Kafka CDC
Consume change events from Kafka and write them into one Paimon table or route them to separate tables in a Paimon database.
See CDC Action Configuration for command syntax, computed columns, and job settings, and Schema Evolution and Type Mapping for shared behavior.
Prerequisites
Place the connector jar matching your Flink runtime in <FLINK_HOME>/lib/:
flink-sql-connector-kafka-*.jar
Supported Formats
Set the event format with --kafka_conf value.format=<format>. Paimon parses insert, update,
and delete events from the message value. Plain JSON represents insert-only data.
value.format | Input |
|---|---|
canal-json | Canal JSON change events |
debezium-json | Debezium JSON change events; include the schema field for source types |
debezium-avro | Debezium Avro using a Confluent schema registry; set schema.registry.url |
maxwell-json | Maxwell JSON change events |
ogg-json | Oracle GoldenGate JSON change events |
json | Plain JSON records, treated as inserts |
aws-dms-json | AWS Database Migration Service change events |
debezium-bson | MongoDB events captured by Debezium |
Check the event metadata before choosing a synchronization mode:
- Missing field types are mapped to
STRING. Preserve the Debeziumschemafield when available. - Database synchronization needs source database and table names to route records. Without these, use a table action and name the target explicitly.
- Missing primary keys can produce an append-only target. Supply
--primary_keyswhen appropriate for update/delete events, and ensure the required key fields are present in those events.
Synchronizing Tables
By using KafkaSyncTableAction in a Flink DataStream job or directly through flink run, users can synchronize one or multiple tables from Kafka's one topic into one Paimon table.
If the target does not exist, the action creates it from the first usable data record found during schema discovery. It then evolves the schema as records arrive. For an existing target, it checks schema compatibility before starting ingestion when source metadata is available.
Example: ingest CDC events
This example assumes the source event contains id and create_time, and derives the pt partition:
<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
kafka_sync_table \
--warehouse hdfs:///path/to/warehouse \
--database test_db \
--table test_table \
--partition_keys pt \
--primary_keys id,pt \
--computed_column 'pt=date_format(create_time,yyyy-MM-dd)' \
--kafka_conf properties.bootstrap.servers=127.0.0.1:9020 \
--kafka_conf topic=order \
--kafka_conf properties.group.id=123456 \
--kafka_conf value.format=canal-json \
--catalog_conf metastore=hive \
--catalog_conf uri=thrift://hive-metastore:9083 \
--table_conf bucket=4 \
--table_conf changelog-producer=input \
--table_conf sink.parallelism=4
Start with an Empty Topic
If schema discovery cannot find a usable data record, create the target table first. Follow Starting with an Empty Topic, then submit the table action with its source connection options.
Example: ingest append-only JSON
Use value.format=json for records such as logs that contain only inserts. This example derives
pt from the source field event_tm and creates a table without primary keys.
<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
kafka_sync_table \
--warehouse hdfs:///path/to/warehouse \
--database test_db \
--table test_table \
--partition_keys pt \
--computed_column 'pt=date_format(event_tm, yyyyMMdd)' \
--kafka_conf properties.bootstrap.servers=127.0.0.1:9020 \
--kafka_conf topic=test_log \
--kafka_conf properties.group.id=123456 \
--kafka_conf value.format=json \
--catalog_conf metastore=hive \
--catalog_conf uri=thrift://hive-metastore:9083 \
--table_conf sink.parallelism=4
Table Action Reference
Command syntax (square brackets indicate optional arguments):
<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
kafka_sync_table \
--warehouse <warehouse-path> \
--database <database-name> \
--table <table-name> \
[--partition_keys <partition_keys>] \
[--primary_keys <primary-keys>] \
[--type_mapping to-string] \
[--computed_column <'column-name=expr-name(args[, ...])'> [--computed_column ...]] \
[--metadata_column <metadata-column> [--metadata_column ...]] \
[--metadata_column_prefix <metadata-column-prefix>] \
[--kafka_conf <kafka-source-conf> [--kafka_conf <kafka-source-conf> ...]] \
[--catalog_conf <paimon-catalog-conf> [--catalog_conf <paimon-catalog-conf> ...]] \
[--table_conf <paimon-table-sink-conf> [--table_conf <paimon-table-sink-conf> ...]]
| Configuration | Description |
|---|---|
--warehouse |
The path to Paimon warehouse. |
--database |
The database name in Paimon catalog. |
--table |
The Paimon table name. |
--partition_keys |
The partition keys for Paimon table. If there are multiple partition keys, connect them with comma, for example "dt,hh,mm". |
--primary_keys |
The primary keys for Paimon table. If there are multiple primary keys, connect them with comma, for example "buyer_id,seller_id". |
--type_mapping |
It is used to specify how to map MySQL data type to Paimon type. Supported options:
|
--sync_primary_keys_from_source_schema |
This is used to specify if primary keys from source should be used in paimon schema if primary keys using --primary_keys are not specified. The default is true. |
--computed_column |
The definitions of computed columns. The argument field is from Kafka topic's table field name. See here for a complete list of configurations. |
--metadata_column |
--metadata_column is used to specify which metadata columns to include in the output schema of the connector. Metadata columns provide additional information related to the source data. Available values are topic, partition, offset, timestamp and timestamp_type (one of 'NoTimestampType', 'CreateTime' or 'LogAppendTime'). |
--metadata_column_prefix |
--metadata_column_prefix is optionally used to set a prefix for metadata columns in the Paimon table to avoid conflicts with existing attributes. For example, with prefix "__kafka_", the metadata column "topic" will be stored as "__kafka_topic" field. |
--kafka_conf |
The configuration for Flink Kafka sources. Each configuration should be specified in the format key=value. properties.bootstrap.servers, topic/topic-pattern, properties.group.id, and value.format are required configurations, others are optional.See its document for a complete list of configurations. |
--catalog_conf |
The configuration for Paimon catalog. Each configuration should be specified in the format "key=value". See here for a complete list of catalog configurations. |
--table_conf |
The configuration for Paimon table sink. Each configuration should be specified in the format "key=value". See here for a complete list of table configurations. |
Synchronizing Databases
By using KafkaSyncDatabaseAction in a Flink DataStream job or directly through flink run, users can route source tables from one or more topics into a Paimon database.
The action uses a shared sink for the selected tables. It creates each missing target from that source table's records and applies supported schema changes to existing targets. Source tables without primary keys are also supported; choose target keys that match the event semantics.
Example: route tables from one topic
<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
kafka_sync_database \
--warehouse hdfs:///path/to/warehouse \
--database test_db \
--kafka_conf properties.bootstrap.servers=127.0.0.1:9020 \
--kafka_conf topic=order \
--kafka_conf properties.group.id=123456 \
--kafka_conf value.format=canal-json \
--catalog_conf metastore=hive \
--catalog_conf uri=thrift://hive-metastore:9083 \
--table_conf bucket=4 \
--table_conf changelog-producer=input \
--table_conf sink.parallelism=4
Example: route tables from multiple topics
Separate topic names with semicolons and quote the argument so the shell passes it as one value.
<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
kafka_sync_database \
--warehouse hdfs:///path/to/warehouse \
--database test_db \
--kafka_conf properties.bootstrap.servers=127.0.0.1:9020 \
--kafka_conf 'topic=order;logistic_order;user' \
--kafka_conf properties.group.id=123456 \
--kafka_conf value.format=canal-json \
--catalog_conf metastore=hive \
--catalog_conf uri=thrift://hive-metastore:9083 \
--table_conf bucket=4 \
--table_conf changelog-producer=input \
--table_conf sink.parallelism=4
Database Action Reference
Command syntax (square brackets indicate optional arguments):
<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
kafka_sync_database \
--warehouse <warehouse-path> \
--database <database-name> \
[--table_mapping <table-name>=<paimon-table-name1> [--table_mapping <table-name2>=<paimon-table-name2> ...]] \
[--table_prefix <paimon-table-prefix>] \
[--table_suffix <paimon-table-suffix>] \
[--table_prefix_db <db-name1>=<table-prefix1> [--table_prefix_db <db-name2>=<table-prefix2> ...]] \
[--table_suffix_db <db-name1>=<table-suffix1> [--table_suffix_db <db-name2>=<table-suffix2> ...]] \
[--including_tables <table-name|name-regular-expr>] \
[--excluding_tables <table-name|name-regular-expr>] \
[--including_dbs <database-name|name-regular-expr>] \
[--excluding_dbs <database-name|name-regular-expr>] \
[--type_mapping to-string] \
[--partition_keys <partition_keys>] \
[--primary_keys <primary-keys>] \
[--computed_column <'column-name=expr-name(args[, ...])'> [--computed_column ...]] \
[--metadata_column <metadata-column> [--metadata_column ...]] \
[--metadata_column_prefix <metadata-column-prefix>] \
[--kafka_conf <kafka-source-conf> [--kafka_conf <kafka-source-conf> ...]] \
[--catalog_conf <paimon-catalog-conf> [--catalog_conf <paimon-catalog-conf> ...]] \
[--table_conf <paimon-table-sink-conf> [--table_conf <paimon-table-sink-conf> ...]]
| Configuration | Description |
|---|---|
--warehouse |
The path to Paimon warehouse. |
--database |
The database name in Paimon catalog. |
--table_mapping |
The table name mapping between source database and Paimon. For example, if you want to synchronize a source table named "test" to a Paimon table named "paimon_test", you can specify "--table_mapping test=paimon_test". Multiple mappings could be specified with multiple "--table_mapping" options. "--table_mapping" has higher priority than "--table_prefix" and "--table_suffix". |
--table_prefix |
The prefix of all Paimon tables to be synchronized except those specified by "--table_mapping" or "--table_prefix_db". For example, if you want all synchronized tables to have "ods_" as prefix, you can specify "--table_prefix ods_". |
--table_suffix |
The suffix of all Paimon tables to be synchronized except those specified by "--table_mapping" or "--table_suffix_db". The usage is same as "--table_prefix". |
--table_prefix_db |
The prefix of the Paimon tables to be synchronized from the specified db. For example, if you want to prefix the tables from db1 with "ods_db1_", you can specify "--table_prefix_db db1=ods_db1_". Multiple mappings could be specified multiple "--table_prefix_db" options. "--table_prefix_db" has higher priority than "--table_prefix". |
--table_suffix_db |
The suffix of the Paimon tables to be synchronized from the specified db. The usage is same as "--table_prefix_db". |
--including_tables |
It is used to specify which source tables are to be synchronized. You must use '|' to separate multiple tables.Quote the expression because '|' is a shell operator, for example: 'a|b|c'. Regular expression is supported, for example, specifying "--including_tables test|paimon.*" means to synchronize table 'test' and all tables start with 'paimon'. |
--excluding_tables |
It is used to specify which source tables are not to be synchronized. The usage is same as "--including_tables". "--excluding_tables" has higher priority than "--including_tables" if you specified both. |
--including_dbs |
It is used to specify the databases within which the tables are to be synchronized. The usage is same as "--including_tables". |
--excluding_dbs |
It is used to specify the databases within which the tables are not to be synchronized. The usage is same as "--excluding_tables". "--excluding_dbs" has higher priority than "--including_dbs" if you specified both. |
--type_mapping |
It is used to specify how to map MySQL data type to Paimon type. Supported options:
|
--computed_column |
The definitions of computed columns. The argument field is from Kafka topic's table field name. See here for a complete list of configurations. NOTICE: It returns null if the referenced column does not exist in the source table. |
--metadata_column |
--metadata_column is used to specify which metadata columns to include in the output schema of the connector. Metadata columns provide additional information related to the source data. Available values are topic, partition, offset, timestamp and timestamp_type (one of 'NoTimestampType', 'CreateTime' or 'LogAppendTime'). |
--metadata_column_prefix |
--metadata_column_prefix is optionally used to set a prefix for metadata columns in the Paimon table to avoid conflicts with existing attributes. For example, with prefix "__kafka_", the metadata column "topic" will be stored as "__kafka_topic" field. |
--eager_init |
It is default false. If true, all relevant tables commiter will be initialized eagerly, which means those tables could be forced to create snapshot. |
--partition_keys |
The partition keys for Paimon table. If there are multiple partition keys, connect them with comma, for example "dt,hh,mm". If the keys are not in source table, the sink table won't set partition keys. |
--multiple_table_partition_keys |
The partition keys for each different Paimon table. If there are multiple partition keys, connect them with comma, for example
|
--primary_keys |
The primary keys for Paimon table. If there are multiple primary keys, connect them with comma, for example "buyer_id,seller_id". If the keys are not provided, but the source has primary keys, the sink table will use source's primary keys. Otherwise, the sink table won't set primary keys. If the keys are not provided, but the source has primary keys, and you don't want to use source's primary keys, use --sync_primary_keys_from_source_schema. |
--sync_primary_keys_from_source_schema |
This is used to specify if primary keys from source should be used in paimon schema if primary keys using --primary_keys are not specified. The default is true. |
--kafka_conf |
The configuration for Flink Kafka sources. Each configuration should be specified in the format key=value. properties.bootstrap.servers, topic/topic-pattern, properties.group.id, and value.format are required configurations, others are optional.See its document for a complete list of configurations. |
--catalog_conf |
The configuration for Paimon catalog. Each configuration should be specified in the format "key=value". See here for a complete list of catalog configurations. |
--table_conf |
The configuration for Paimon table sink. Each configuration should be specified in the format "key=value". See here for a complete list of table configurations. |
Additional Kafka Configuration
Pass this option through --kafka_conf in addition to the source connection settings:
| Option | Default | Description |
|---|---|---|
schema.registry.url | None | Confluent schema registry URL, required for value.format=debezium-avro. |
Debezium-bson
For MongoDB events captured by Debezium, use value.format=debezium-bson. See
Kafka Debezium BSON Format for full-document requirements, the Kafka key
fallback, and BSON-to-string conversion examples.