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.
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
| Scope | Source settings |
|---|---|
| One table | Set database and table. |
| One database | Set database; omit table. |
| Whole warehouse | Omit 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-intervalhas 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 type | Flink CDC type |
|---|---|
TINYINT, SMALLINT, INT, BIGINT | Corresponding integer type |
FLOAT, DOUBLE | Corresponding floating-point type |
DECIMAL(p, s) | DECIMAL(p, s) |
BOOLEAN | BOOLEAN |
DATE | DATE |
TIMESTAMP | TIMESTAMP |
TIMESTAMP_LTZ | TIMESTAMP_LTZ |
CHAR(n), VARCHAR(n) | Corresponding character type |