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'
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
;
