DataZen Documentation
DataZen User Guide

PUSH MESSAGE

Overview

The Push Message operation allows you to send the Upsert and Delete streams to a Messaging Bus endpoint. Each messaging platform offers different connection and processing options.

In some scenarios, you may need to send a single item at a time; use the BATCH option to control this behavior.

OPTIONS Parameter

The OPTIONS clause on ON_UPSERT and ON_DELETE configures messaging-platform settings as a comma-separated list of key-value pairs. Supported keys vary by messaging engine.

  • Option names are case-insensitive.
  • Values are trimmed; surrounding " or ' characters are stripped automatically.
  • DataZen functions and variables are supported in option values.

Example for Kafka:

OPTIONS 'ApplicationId=testapp,CorrelationId=@jobguid'
Kafka
Option Description
TopicName Target Kafka topic name.
KeyFieldSeparator Separator used when multiple upsert key columns are combined into the Kafka message key. The key column list itself comes from the pipeline key settings, not from OPTIONS.
ApplicationId Application identifier associated with the message.
CorrelationId Correlation identifier for relating related messages.
ReplyTo Reply-to destination for response routing.
Azure Service Bus
Option Description
QueueOrTopicName Target queue or topic name.
AzBusQueueKind Queue or topic kind. Supported values: Dynamic (default), Queue, or Topic (also accepts 0, 1, or 2).
ForwardProperties When true, 1, yes, or y, forwards additional message properties.
MessageId Message identifier.
Subject Message subject / label.
To Destination address for the message.
ReplyTo Reply-to address.
CorrelationId Correlation identifier for relating related messages.
Azure Event Hub
Option Description
EventHubName Target Event Hub name.
MessageId Message identifier.
CorrelationId Correlation identifier for relating related messages.
AWS SQS
Option Description
QueueName Target SQS queue name.
MessageGroupId FIFO queues only. Message group ID. Can be a static value or a {{field}} reference. When blank, no group ID is sent.
DeduplicationIdField FIFO queues only. Field used to build a deduplication ID. When blank, no deduplication ID is sent and the queue is assumed to use content-based deduplication.
RabbitMQ
Option Description
RoutingKey Routing key used when publishing the message.
ReplyTo Reply-to queue or address.
Priority Message priority (integer).
MessageId Message identifier.
AppId Application identifier associated with the message.
Expiration Message expiration in milliseconds (long integer).
Persist When true, 1, yes, or y, marks the message as persistent.
Google Pub/Sub
Option Description
TopicName Target Pub/Sub topic name.
MSMQ
Option Description
QueueName Target MSMQ queue name.
Label Message label.
Formatter Message formatter used when writing to the queue.
CorrelationId Correlation identifier for relating related messages.
IsAuthenticated When true, 1, yes, or y, marks the message as authenticated.

Syntax

Sends a payload to a messaging platform for upserted and/or deleted records one at a time or in a batch.

PUSH INTO MESSAGE [CONNECTION]
    {ON_UPSERT
        PAYLOAD (...) 
        OPTIONS (...)
    }
    {ON_DELETE
        PAYLOAD (...) 
        OPTIONS (...)
    }
    { WITH
        { TIMEOUT N }
        { BATCH N }
        { DISCARD_ON_SUCCESS }
        { DISCARD_ON_ERROR }
        { DLQ_ON_ERROR '...' }
        { RETRY < LINEAR | EXPONENTIAL > (N,N) }
    }
;

ON_UPSERT

Section that defines the message operation to execute when processing upserted records

ON_DELETE

Section that defines the message operation to execute when processing deleted records

PAYLOAD

The payload to use or build for upsert or delete operations; may contain DataZen functions and pipeline variables

OPTIONS

The messaging platform-specific set of options for upsert or delete operations as a comma-separated list of key=value settings (ex: applicationId=test)

TIMEOUT

A timeout value in seconds

BATCH

The number of records to process in a single call

RETRY_INTERVAL

The retry interval in seconds (default: 1)

RETRY_COUNT

The maximum number of retries to perform (default: 1)

DISCARD_ON_SUCCESS

Deletes the change log after successful completion of the push operation

DISCARD_ON_ERROR

Deletes the change log if the push operation failed

DLQ_ON_ERROR

Moves the change log to a subfolder or directory if the push operation failed

RETRY EXPONENTIAL

Retries the operation on failure up to N times, waiting P seconds exponentially longer every time (N,P)

RETRY LINEAR

Retries the operation on failure up to N times, waiting P seconds every time (N,P)

Example 1


-- Load the next inserts previously captured in a cloud folder using
-- CAPTURE 'mycdc' INSERT ON KEYS 'guid' WITH PATH [adls2.0] 'container'

LOAD UPSERTS FROM [adls2.0] 'mycdc' PATH '/container' KEY_COLUMNS 'guid';

-- Sends a message for each record to a Kafka endpoint
PUSH INTO MESSAGE [kafka-cloud] 
	ON_UPSERT 
		PAYLOAD ({"Title": "{{title}}")
		OPTIONS (applicationId=integration,topicName=rssfeed)
WITH 
	BATCH 1
;

Example 2


-- Retrieves newly created tables from a database since the last execution
-- @highwatermarknull will absorb surrounding quotes the first time the command runs 
SELECT * FROM DB [sql2017]
(SELECT name, object_id, modify_date FROM sys.tables WHERE modify_date > ISNULL(@highwatermarknull, '')
WITH HWM 'modify_date';

-- Sends a single payload of rss feeds 100 at a time to an AWS SQS endpoint using the @contactjson operation 
PUSH INTO MESSAGE [aws-cloud] 
	ON_UPSERT 
		PAYLOAD (@concatjson({"tableName": "{{name}}", "modifiedOn": "modify_date}}", "object_id": object_id))
		OPTIONS (QueueName=integration)
WITH 
	BATCH 100
;

Example 3


-- Retrieves newly created tables from a database since the last execution
-- @highwatermarknull will absorb surrounding quotes the first time the command runs 
SELECT * FROM DB [sql2017]
(SELECT name, object_id, modify_date FROM sys.tables WHERE modify_date > ISNULL(@highwatermarknull, '')
WITH HWM 'modify_date';

-- Build a JSON document inline for up to 50 records at a time 
-- creating 1 row per batch (100 records will yield 2 rows)
ADD COLUMN 'json' FORMAT JSON;
ZIP COLUMN 'json' FORMAT JSON BATCH 50;

-- Sends a single payload of data 50 at a time to an Azure Bus endpoint using the @contactjson operation 
PUSH INTO MESSAGE [azbus-cloud] 
	ON_UPSERT 
		PAYLOAD ({{json}})
		OPTIONS (QueueOrTopicName=integration,Subject=new or modified tables,CorrelationId=#rndguid())
WITH 
	BATCH 1
;