Skip to main content
Version: Current

Table (Sink)

version: Enterprise Mode: Streaming

Table sink component is automatically generated based on the Flink Table API integration definition. The name of the component will be Integration-name Table. Check the Table API integration for more information on Flink Table API connectors.

Writing a changelog (CDC)​

For a table whose connector can consume a changelog, the sink has an additional CDC Operation parameter.

CDC Operation has to be one of:

OperationMeaning
INSERTInsertion operation
DELETEDeletion operation
UPDATE_BEFOREUpdate operation with the previous content of the updated row
UPDATE_AFTERUpdate operation with new content of the updated row

The effect of each operation depends on the connector used by a table. For example, connectors applying changes by PRIMARY KEY treat INSERT and UPDATE_AFTER as an upsert and ignore UPDATE_BEFORE.

Any expression can be used as the operation. For a plain stream the parameter defaults to 'INSERT'.

Example: a customers table with columns id (primary key) and name, written by the jdbc connector, with the sink Value set to {id: #input.id, name: #input.name}:

CDC OperationValueTable contents after the record
INSERT{id: 1, name: 'Alice'}(1, Alice)
INSERT{id: 2, name: 'Bob'}(1, Alice), (2, Bob)
UPDATE_AFTER{id: 1, name: 'Alice-updated'}(1, Alice-updated), (2, Bob)
DELETE{id: 2, name: 'Bob'}(1, Alice-updated)