Skip to main content

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.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

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 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:
  • "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.
--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:
  • "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.
--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:

OptionDefaultDescription
value.formatNoneEvent format from the supported formats table.
topicNoneSemicolon-separated topic names. Set either topic or topic-pattern.
topic-patternNoneRegular expression for topics to subscribe to.
pulsar.startCursor.fromMessageIdEARLIESTStart at ledgerId,entryId,partitionIndex, EARLIEST, or LATEST. Mutually exclusive with fromPublishTime.
pulsar.startCursor.fromPublishTimeNoneStart at the given message publish timestamp (milliseconds).
pulsar.startCursor.fromMessageIdInclusivetrueInclude the specified message ID. Does not apply to EARLIEST or LATEST.
pulsar.stopCursor.atMessageIdNoneStop before the specified message ID, expressed as ledgerId,entryId,partitionIndex or LATEST.
pulsar.stopCursor.afterMessageIdNoneStop after consuming the specified message ID, expressed as ledgerId,entryId,partitionIndex or LATEST.
pulsar.stopCursor.atEventTimeNoneStop before messages whose event timestamp is greater than or equal to this value (milliseconds).
pulsar.stopCursor.afterEventTimeNoneStop before messages whose event timestamp is greater than this value (milliseconds).
pulsar.source.unboundedtrueRun as an unbounded source.
schema.registry.urlNoneConfluent schema registry URL, required for value.format=debezium-avro.