Skip to main content

MySQL CDC

Synchronize MySQL tables into one Paimon table, or synchronize selected tables into 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 these dependencies in <FLINK_HOME>/lib/ alongside the Paimon Flink bundled jar. The dependency baseline for this Paimon version is Flink CDC 3.5.0.

All source tables selected by the table action must have primary keys, including when you supply --primary_keys for the target. Schemas must be compatible when merging multiple source tables. The partitioned examples below assume each source table has id and create_time fields.

Synchronizing Tables​

By using MySqlSyncTableAction in a Flink DataStream job or directly through flink run, users can synchronize one or multiple tables from MySQL into one Paimon table.

If the Paimon table you specify does not exist, this action will automatically create the table. Its schema will be derived from all specified MySQL tables. If the Paimon table already exists, its schema will be compared against the schema of all specified MySQL tables.

Example 1: synchronize tables into one Paimon table​

<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
mysql_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)' \
--mysql_conf hostname=127.0.0.1 \
--mysql_conf username=root \
--mysql_conf password=123456 \
--mysql_conf database-name='source_db' \
--mysql_conf table-name='source_table1|source_table2' \
--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

As example shows, the mysql_conf's table-name supports regular expressions to monitor multiple tables that satisfy the regular expressions. The schemas of all the tables will be merged into one Paimon table schema.

Example 2: synchronize shards into one Paimon table​

You can also set 'database-name' with a regular expression to capture multiple databases. A typical scenario is that a table 'source_table' is split into database 'source_db1', 'source_db2' ..., then you can synchronize data of all the 'source_table's into one Paimon table.

<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
mysql_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)' \
--mysql_conf hostname=127.0.0.1 \
--mysql_conf username=root \
--mysql_conf password=123456 \
--mysql_conf database-name='source_db.+' \
--mysql_conf table-name='source_table' \
--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

Table Action Reference​

Command syntax (square brackets indicate optional arguments):

<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
mysql_sync_table \
--warehouse <warehouse-path> \
--database <database-name> \
--table <table-name> \
[--partition_keys <partition_keys>] \
[--primary_keys <primary-keys>] \
[--type_mapping <option1,option2...>] \
[--computed_column <'column-name=expr-name(args[, ...])'> [--computed_column ...]] \
[--metadata_column <metadata-column>] \
[--mysql_conf <mysql-cdc-source-conf> [--mysql_conf <mysql-cdc-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 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 MySQL table field name. See here for a complete list of configurations.
--metadata_column
--metadata_column is used to specify which metadata columns to include in the output schema of the connector. Metadata columns provide additional information related to the source data, for example: --metadata_column table_name,database_name,op_ts. See its document for a complete list of available metadata.
--mysql_conf
The configuration for Flink CDC MySQL sources. Each configuration should be specified in the format "key=value". hostname, username, password, database-name and table-name 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 MySqlSyncDatabaseAction in a Flink DataStream job or directly through flink run, users can synchronize the whole MySQL database into one Paimon database.

Only tables with primary keys will be synchronized.

The action creates a target for each selected source table and checks existing target schemas for compatibility. When shard merging is enabled, only source tables with the same table name are merged into the same target.

Choose a Synchronization Mode​

ModeBehavior
divided (default)Build a separate sink for each table selected at startup. Restore the job with an expanded table filter to add previously excluded tables.
combinedUse a shared sink and support newly created source tables that match the filters without restarting.

Adding a previously excluded table with historical data is different from capturing a table created after the job starts. Use the savepoint workflow below when expanding the selected set.

Per-table table configuration​

Use repeated --table_conf_by_table arguments when different MySQL source tables need different Paimon table properties:

--table_conf bucket=4 \
--table_conf changelog-producer=input \
--table_conf_by_table orders:bucket=8 \
--table_conf_by_table users:bucket=2

The global --table_conf is the default. A matching per-table option overrides the same key; users above uses bucket=2, while an unconfigured table uses bucket=4. Repeat the option for each property. The source name is the MySQL table name, not the generated Paimon table name.

This option is supported for both divided and combined database synchronization. It applies only to Paimon table properties. Sink and job options such as sink.parallelism, writer resources, and committer resources are shared runtime settings and must be configured globally with --table_conf; they cannot be overridden per table.

Example 1: synchronize entire database​

<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
mysql_sync_database \
--warehouse hdfs:///path/to/warehouse \
--database test_db \
--mysql_conf hostname=127.0.0.1 \
--mysql_conf username=root \
--mysql_conf password=123456 \
--mysql_conf database-name=source_db \
--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 2: synchronize newly added tables under database​

Let's say at first a Flink job is synchronizing tables [product, user, address] under database source_db. The command to submit the job looks like:

<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
mysql_sync_database \
--warehouse hdfs:///path/to/warehouse \
--database test_db \
--mysql_conf hostname=127.0.0.1 \
--mysql_conf username=root \
--mysql_conf password=123456 \
--mysql_conf database-name=source_db \
--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 \
--including_tables 'product|user|address'

At a later point we would like the job to also synchronize tables [order, custom], which contains history data. We can achieve this by recovering from the previous savepoint of the job and thus reusing existing state of the job. The recovered job will first snapshot newly added tables, and then continue reading changelog from previous position automatically.

The command to restore from a savepoint and add new tables to synchronize looks like:

<FLINK_HOME>/bin/flink run \
--fromSavepoint savepointPath \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
mysql_sync_database \
--warehouse hdfs:///path/to/warehouse \
--database test_db \
--mysql_conf hostname=127.0.0.1 \
--mysql_conf username=root \
--mysql_conf password=123456 \
--mysql_conf database-name=source_db \
--catalog_conf metastore=hive \
--catalog_conf uri=thrift://hive-metastore:9083 \
--table_conf bucket=4 \
--including_tables 'product|user|address|order|custom'
info

You can set --mode combined to enable synchronizing newly added tables without restarting job.

Example 3: synchronize and merge multiple shards​

Let's say you have multiple database shards db1, db2, ... and each database has tables tbl1, tbl2, .... You can synchronize all the db.+.tbl.+ into tables test_db.tbl1, test_db.tbl2 ... by following command:

<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
mysql_sync_database \
--warehouse hdfs:///path/to/warehouse \
--database test_db \
--mysql_conf hostname=127.0.0.1 \
--mysql_conf username=root \
--mysql_conf password=123456 \
--mysql_conf database-name='db.+' \
--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 \
--including_tables 'tbl.+'

By setting database-name to a regular expression, the synchronization job will capture all tables under matched databases and merge tables of the same name into one table.

info

You can set --merge_shards false to prevent merging shards. The synchronized tables will be named to 'databaseName_tableName' to avoid potential name conflict.

Database Action Reference​

Command syntax (square brackets indicate optional arguments):

<FLINK_HOME>/bin/flink run \
/path/to/paimon-flink-action-2.2-SNAPSHOT.jar \
mysql_sync_database \
--warehouse <warehouse-path> \
--database <database-name> \
[--ignore_incompatible <true/false>] \
[--merge_shards <true/false>] \
[--table_prefix <paimon-table-prefix>] \
[--table_suffix <paimon-table-suffix>] \
[--including_tables <mysql-table-name|name-regular-expr>] \
[--excluding_tables <mysql-table-name|name-regular-expr>] \
[--mode <sync-mode>] \
[--metadata_column <metadata-column>] \
[--type_mapping <option1,option2...>] \
[--partition_keys <partition_keys>] \
[--primary_keys <primary-keys>] \
[--mysql_conf <mysql-cdc-source-conf> [--mysql_conf <mysql-cdc-source-conf> ...]] \
[--catalog_conf <paimon-catalog-conf> [--catalog_conf <paimon-catalog-conf> ...]] \
[--table_conf <paimon-table-sink-conf> [--table_conf <paimon-table-sink-conf> ...]] \
[--table_conf_by_table <source-table>:<key>=<value> [--table_conf_by_table <source-table>:<key>=<value> ...]]
Configuration Description
--warehouse
The path to Paimon warehouse.
--database
The database name in Paimon catalog.
--ignore_incompatible
It is default false, in this case, if MySQL table name exists in Paimon and their schema is incompatible,an exception will be thrown. You can specify it to true explicitly to ignore the incompatible tables and exception.
--merge_shards
It is default true, in this case, if some tables in different databases have the same name, their schemas will be merged and their records will be synchronized into one Paimon table. Otherwise, each table's records will be synchronized to a corresponding Paimon table, and the Paimon table will be named to 'databaseName_tableName' to avoid potential name conflict.
--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.
--mode
It is used to specify synchronization mode.
Possible values:
  • "divided" (the default mode if you haven't specified one): start a sink for each table, the synchronization of the new table requires restarting the job.
  • "combined": start a single combined sink for all tables, the new table will be automatically synchronized.
--metadata_column
--metadata_column is used to specify which metadata columns to include in the output schema of the connector. Metadata columns provide additional information related to the source data, for example: --metadata_column table_name,database_name,op_ts. See its document for a complete list of available metadata.
--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 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.
--mysql_conf
The configuration for Flink CDC MySQL sources. Each configuration should be specified in the format "key=value". Set hostname, username, password, and database-name. Select tables with --including_tables and --excluding_tables; do not set table-name for the database action. 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.
--table_conf_by_table
The configuration for one Paimon table sink, overriding the same key from --table_conf for that table. Each configuration should be specified in the format "source_table:key=value", where source_table is the MySQL table name. Only table properties can be overridden per table; options prefixed with "sink." must be set with --table_conf.

When submitting a MySQL CDC action with Apache Flink Kubernetes Operator, put the Paimon action jar in a Flink image and pass the same action arguments through job.args. Do not run <FLINK_HOME>/bin/flink run inside the container. The operator submits the jar from job.jarURI.

Build a Flink image which contains the Paimon bundled jar, the MySQL CDC dependencies, and the Paimon action jar. For example:

# Use a Flink base image which matches your Flink Kubernetes Operator and Paimon bundled jar version.
FROM flink:<flink-version>

COPY jars/paimon-flink-1.20-2.2-SNAPSHOT.jar /opt/flink/lib/
COPY jars/flink-sql-connector-mysql-cdc-3.5.0.jar /opt/flink/lib/
COPY jars/mysql-connector-java-8.0.27.jar /opt/flink/lib/
COPY jars/paimon-flink-action-2.2-SNAPSHOT.jar /opt/flink/usrlib/

The following FlinkDeployment starts a mysql_sync_database action. Use mysql_sync_table and add the --table argument if you want to synchronize into one Paimon table.

apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
name: paimon-mysql-cdc
spec:
image: <registry>/flink-paimon-mysql-cdc:latest
flinkVersion: v1_20
serviceAccount: flink
flinkConfiguration:
taskmanager.numberOfTaskSlots: "4"
state.checkpoints.dir: hdfs:///path/to/checkpoints
state.savepoints.dir: hdfs:///path/to/savepoints
jobManager:
resource:
memory: "2048m"
cpu: 1
taskManager:
resource:
memory: "4096m"
cpu: 2
job:
jarURI: local:///opt/flink/usrlib/paimon-flink-action-2.2-SNAPSHOT.jar
entryClass: org.apache.paimon.flink.action.FlinkActions
args:
- mysql_sync_database
- --warehouse
- hdfs:///path/to/warehouse
- --database
- test_db
- --mysql_conf
- hostname=mysql.default.svc.cluster.local
- --mysql_conf
- username=root
- --mysql_conf
- password=123456
- --mysql_conf
- database-name=source_db
- --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
parallelism: 4
upgradeMode: savepoint

If you use Hive catalog or a Hadoop-compatible filesystem, also make sure the image or pod mounts contain the required Hadoop, Hive, and filesystem configuration files. In production, pass sensitive values such as the MySQL password through Kubernetes Secrets instead of writing them directly in the manifest.

FAQ​

Garbled Characters​

Set env.java.opts.all: -Dfile.encoding=UTF-8 in the Flink configuration (config.yaml for Flink 1.19 and later, flink-conf.yaml for earlier versions). Before Flink 1.17, use env.java.opts.

Table and Column Comments​

  • For comments on newly created MySQL tables, set --mysql_conf jdbc.properties.useInformationSchema=true.
  • For comment changes, set --mysql_conf debezium.include.schema.comments=true.