Pulsar CDC
Consume change events from Pulsar 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-connector-pulsar-*.jar
Supported Formats
Set the event format with --pulsar_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 |
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 PulsarSyncTableAction in a Flink DataStream job or directly through flink run, users can synchronize one or multiple tables from Pulsar'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 \
pulsar_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)' \
--pulsar_conf topic=order \
--pulsar_conf value.format=canal-json \
--pulsar_conf pulsar.client.serviceUrl=pulsar://127.0.0.1:6650 \
--pulsar_conf pulsar.admin.adminUrl=http://127.0.0.1:8080 \
--pulsar_conf pulsar.consumer.subscriptionName=paimon-tests \
--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 \
pulsar_sync_table \
--warehouse hdfs:///path/to/warehouse \
--database test_db \
--table test_table \
--partition_keys pt \
--computed_column 'pt=date_format(event_tm, yyyyMMdd)' \
--pulsar_conf pulsar.client.serviceUrl=pulsar://127.0.0.1:6650 \
--pulsar_conf pulsar.admin.adminUrl=http://127.0.0.1:8080 \
--pulsar_conf topic=test_log \
--pulsar_conf pulsar.consumer.subscriptionName=paimon-tests \
--pulsar_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 \
pulsar_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 ...]] \
[--pulsar_conf <pulsar-source-conf> [--pulsar_conf <pulsar-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 Pulsar topic's table field name. See here for a complete list of configurations. |
--pulsar_conf |
The configuration for Flink Pulsar sources. Each configuration should be specified in the format key=value. topic/topic-pattern, value.format, pulsar.client.serviceUrl, pulsar.admin.adminUrl, and pulsar.consumer.subscriptionName 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 PulsarSyncDatabaseAction 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 \
pulsar_sync_database \
--warehouse hdfs:///path/to/warehouse \
--database test_db \
--pulsar_conf topic=order \
--pulsar_conf value.format=canal-json \
--pulsar_conf pulsar.client.serviceUrl=pulsar://127.0.0.1:6650 \
--pulsar_conf pulsar.admin.adminUrl=http://127.0.0.1:8080 \
--pulsar_conf pulsar.consumer.subscriptionName=paimon-tests \
--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 \
pulsar_sync_database \
--warehouse hdfs:///path/to/warehouse \
--database test_db \
--pulsar_conf 'topic=order;logistic_order;user' \
--pulsar_conf value.format=canal-json \
--pulsar_conf pulsar.client.serviceUrl=pulsar://127.0.0.1:6650 \
--pulsar_conf pulsar.admin.adminUrl=http://127.0.0.1:8080 \
--pulsar_conf pulsar.consumer.subscriptionName=paimon-tests \
--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 \
pulsar_sync_database \
--warehouse <warehouse-path> \
--database <database-name> \
[--table_prefix <paimon-table-prefix>] \
[--table_suffix <paimon-table-suffix>] \
[--including_tables <table-name|name-regular-expr>] \
[--excluding_tables <table-name|name-regular-expr>] \
[--type_mapping to-string] \
[--partition_keys <partition_keys>] \
[--primary_keys <primary-keys>] \
[--pulsar_conf <pulsar-source-conf> [--pulsar_conf <pulsar-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_prefix |
The prefix of all Paimon tables to be synchronized. 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. The usage is same as "--table_prefix". |
--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. |
--type_mapping |
It is used to specify how to map MySQL data type to Paimon type. Supported options:
|
--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. |
--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. |
--pulsar_conf |
The configuration for Flink Pulsar sources. Each configuration should be specified in the format key=value. topic/topic-pattern, value.format, pulsar.client.serviceUrl, pulsar.admin.adminUrl, and pulsar.consumer.subscriptionName 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 Pulsar Configuration
Pass these options through --pulsar_conf:
| Option | Default | Description |
|---|---|---|
value.format | None | Event format from the supported formats table. |
topic | None | Semicolon-separated topic names. Set either topic or topic-pattern. |
topic-pattern | None | Regular expression for topics to subscribe to. |
pulsar.startCursor.fromMessageId | EARLIEST | Start at ledgerId,entryId,partitionIndex, EARLIEST, or LATEST. Mutually exclusive with fromPublishTime. |
pulsar.startCursor.fromPublishTime | None | Start at the given message publish timestamp (milliseconds). |
pulsar.startCursor.fromMessageIdInclusive | true | Include the specified message ID. Does not apply to EARLIEST or LATEST. |
pulsar.stopCursor.atMessageId | None | Stop before the specified message ID, expressed as ledgerId,entryId,partitionIndex or LATEST. |
pulsar.stopCursor.afterMessageId | None | Stop after consuming the specified message ID, expressed as ledgerId,entryId,partitionIndex or LATEST. |
pulsar.stopCursor.atEventTime | None | Stop before messages whose event timestamp is greater than or equal to this value (milliseconds). |
pulsar.stopCursor.afterEventTime | None | Stop before messages whose event timestamp is greater than this value (milliseconds). |
pulsar.source.unbounded | true | Run as an unbounded source. |
schema.registry.url | None | Confluent schema registry URL, required for value.format=debezium-avro. |