Skip to main content

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.formatInput
canal-jsonCanal JSON change events
debezium-jsonDebezium JSON change events; include the schema field for source types
debezium-avroDebezium Avro using a Confluent schema registry; set schema.registry.url
maxwell-jsonMaxwell JSON change events
ogg-jsonOracle GoldenGate JSON change events
jsonPlain JSON records, treated as inserts
aws-dms-jsonAWS Database Migration Service change events
debezium-bsonMongoDB events captured by Debezium

Check the event metadata before choosing a synchronization mode:

  • Missing field types are mapped to STRING. Preserve the Debezium schema field 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_keys when 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:
  • "tinyint1-not-bool": maps MySQL TINYINT(1) to TINYINT instead of BOOLEAN.
  • "to-nullable": ignores all NOT NULL constraints (except for primary keys). This is used to solve the problem that Flink cannot accept the MySQL 'ALTER TABLE ADD COLUMN column type NOT NULL DEFAULT x' operation.
  • "to-string": maps all MySQL types to STRING.
  • "char-to-string": maps MySQL CHAR(length)/VARCHAR(length) types to STRING.
  • "longtext-to-bytes": maps MySQL LONGTEXT types to BYTES.
  • "bigint-unsigned-to-bigint": maps MySQL BIGINT UNSIGNED, BIGINT UNSIGNED ZEROFILL, SERIAL to BIGINT. You should ensure overflow won't occur when using this option.
  • "decimal-no-change": Ignore decimal type change.
--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:
  • "tinyint1-not-bool": maps MySQL TINYINT(1) to TINYINT instead of BOOLEAN.
  • "to-nullable": ignores all NOT NULL constraints (except for primary keys). This is used to solve the problem that Flink cannot accept the MySQL 'ALTER TABLE ADD COLUMN column type NOT NULL DEFAULT x' operation.
  • "to-string": maps all MySQL types to STRING.
  • "char-to-string": maps MySQL CHAR(length)/VARCHAR(length) types to STRING.
  • "longtext-to-bytes": maps MySQL LONGTEXT types to BYTES.
  • "bigint-unsigned-to-bigint": maps MySQL BIGINT UNSIGNED, BIGINT UNSIGNED ZEROFILL, SERIAL to BIGINT. You should ensure overflow won't occur when using this option.
  • "decimal-no-change": Ignore decimal type change.
--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
  • --multiple_table_partition_keys tableName1=col1,col2.col3
  • --multiple_table_partition_keys tableName2=col4,col5.col6
  • --multiple_table_partition_keys tableName3=col7,col8.col9
  • If the keys are not in source table, the sink table won't set partition keys.
    --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:

    OptionDefaultDescription
    schema.registry.urlNoneConfluent 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.