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
| Source | Synchronize into one Paimon table | Synchronize into a Paimon database |
|---|---|---|
| MySQL | One source table, or merge multiple tables and shards | Synchronize selected tables; optionally merge same-name shards |
| PostgreSQL | One source table, or merge tables across schemas in one database | No database action |
| MongoDB | One collection | Synchronize selected collections |
| Kafka | Change events from one topic | Route tables from one or more topics |
| Pulsar | Change events from one topic | Route tables from one or more topics |
| Flink DataStream API | Ingest custom CDC records | See the API guide for multi-table sinks |
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 YAML Pipelines
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
- Set up Paimon on Flink and install the source connector dependencies
listed in the source guide. Use the
paimon-flink-action-2.2-SNAPSHOT.jarfor the job. - Choose the target warehouse, database, table layout, and primary keys. In action commands,
--databaseand--tableidentify the Paimon target; source names belong in--mysql_conf,--postgres_conf,--mongodb_conf,--kafka_conf, or--pulsar_conf. - 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.
- 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.
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.
- Job settings: configure checkpointing, the job name, catalog properties, and table properties.
- Starting with an empty topic: create the target schema before running a Kafka or Pulsar table action.