Skip to main content

Paimon as a Flink CDC Source

Use the Paimon pipeline source connector to read Paimon tables and send changes to an external system through a Flink CDC YAML pipeline.

Paimon tables feed a Flink CDC pipeline, which sends row and schema changes to an external sink.

To write into Paimon using a YAML pipeline, see Flink CDC's Paimon sink connector. For Paimon's source-specific ingestion commands, return to the CDC ingestion overview.

Capabilities​

  • Read a Paimon warehouse, database, or table.
  • Emit row events and schema changes to the pipeline.
  • Discover newly created tables when the configured scope includes multiple tables.

Create a Pipeline​

This example reads default.test_table from a filesystem catalog and writes to Doris. Make the Paimon source and destination connector dependencies available to the Flink CDC runtime, and replace the catalog path and sink connection values for your deployment.

source:
type: paimon
name: Paimon Source
database: default
table: test_table
catalog.properties.metastore: filesystem
catalog.properties.warehouse: /path/warehouse

sink:
type: doris
name: Doris Sink
fenodes: 127.0.0.1:8030
username: root
password: pass

pipeline:
name: Paimon to Doris Pipeline
parallelism: 2

Pipeline Connector Options​

ScopeSource settings
One tableSet database and table.
One databaseSet database; omit table.
Whole warehouseOmit both database and table.

Use table.discovery-interval to control discovery of new tables when the scope is not fixed to one table. This is separate from waiting for the next snapshot of an existing table.

Key Default Type Description
database
(none) String Name of the database to be scanned. By default, all databases will be scanned.
table
(none) String Name of the table to be scanned. By default, all tables will be scanned.
table.discovery-interval
1 min Duration The discovery interval of new tables. Only effective when database or table is not set.

Catalog Options​

Prefix catalog settings with catalog.properties.. For example, catalog.properties.warehouse becomes the catalog's warehouse option, and catalog.properties.metastore becomes metastore. See Configurations for catalog options.

Usage Notes​

  • Data updates for primary key tables (-U, +U) will be replaced with -D and +I.
  • Does not support dropping tables. If you need to drop a table from the Paimon warehouse, please restart the Flink CDC job after performing the drop operation. When the job restarts, it will stop reading data from the dropped table, and the target table in the external system will remain unchanged from its state before the job was stopped.
  • Data from the same table will be consumed by the same Flink source subtask. If the amount of data varies significantly across different tables, performance bottlenecks caused by data skew may be observed in Flink CDC jobs.
  • If the CDC job has consumed up to the latest snapshot of a table and the next snapshot is not available yet, the monitoring and consumption of this table may be temporarily paused until continuous.discovery-interval has passed.

Data Type Mapping​

Common scalar mappings are shown below. The connector converts through Flink's logical types; acceptance of a type also depends on the destination connector.

Paimon typeFlink CDC type
TINYINT, SMALLINT, INT, BIGINTCorresponding integer type
FLOAT, DOUBLECorresponding floating-point type
DECIMAL(p, s)DECIMAL(p, s)
BOOLEANBOOLEAN
DATEDATE
TIMESTAMPTIMESTAMP
TIMESTAMP_LTZTIMESTAMP_LTZ
CHAR(n), VARCHAR(n)Corresponding character type