Skip to main content

Overview

Use Paimon CDC actions to ingest database changes or change events from a message queue into Paimon tables. The actions create target tables and apply supported schema changes while the Flink job is running.

Choose an Ingestion Path​

SourceSynchronize into one Paimon tableSynchronize into a Paimon database
MySQLOne source table, or merge multiple tables and shardsSynchronize selected tables; optionally merge same-name shards
PostgreSQLOne source table, or merge tables across schemas in one databaseNo database action
MongoDBOne collectionSynchronize selected collections
KafkaChange events from one topicRoute tables from one or more topics
PulsarChange events from one topicRoute tables from one or more topics
Flink DataStream APIIngest custom CDC recordsSee the API guide for multi-table sinks

Table synchronization merges selected sources into one target; database synchronization routes each source table to its target.

Choose table synchronization when you want one target table. If several source tables feed that target, their schemas must be compatible and their primary keys must identify rows across all sources. Rows with the same target primary key represent the same record.

Choose database synchronization when you want to preserve separate tables. Table filters, naming rules, primary-key requirements, and discovery of new tables depend on the source action; check its guide before starting the job.

Flink CDC also supports pipelines described in YAML. To write into Paimon through that framework, see the Flink CDC Paimon sink connector. To read from Paimon and send changes to another system, see Paimon as a Flink CDC Source. The action options in this section, such as --table_conf, belong to the Paimon action CLI.

Before You Start​

  1. Set up Paimon on Flink and install the source connector dependencies listed in the source guide. Use the paimon-flink-action-2.2-SNAPSHOT.jar for the job.
  2. Choose the target warehouse, database, table layout, and primary keys. In action commands, --database and --table identify the Paimon target; source names belong in --mysql_conf, --postgres_conf, --mongodb_conf, --kafka_conf, or --pulsar_conf.
  3. Check the source's schema and change-event requirements. For Kafka and Pulsar, the event format determines which types, table names, and keys can be recovered.
  4. Start with a source-specific example, then add computed columns and job settings.

What Is Schema Evolution​

An ordinary Flink SQL INSERT INTO ... SELECT ... job uses the schemas declared for that job. Adding a column in the source does not automatically add it to the SQL definitions and target table. Paimon CDC actions can carry the new field and its values through to Paimon without restarting for a supported column addition.

A declared Flink SQL schema keeps the original three fields, while a Paimon CDC action adds field_4 to the target.

Supported Changes​

Schema evolution is not full DDL replication. Column additions and compatible type changes can be applied; dropping or renaming source objects does not perform the equivalent operation on existing Paimon objects. See Schema Evolution and Type Mapping for the supported changes, limitations, and mapping options.

Configuration Reference​

  • Computed columns: derive partition values and other fields, with temporal and string functions.
  • Type mapping: control how source types become Paimon types, including MySQL-specific mappings.