Table (Sink)
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:
| Operation | Meaning |
|---|---|
INSERT | Insertion operation |
DELETE | Deletion operation |
UPDATE_BEFORE | Update operation with the previous content of the updated row |
UPDATE_AFTER | Update 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 Operation | Value | Table 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) |