This topic describes the schema evolution configurations for Flink CDC data ingestion jobs.
Schema evolution configurations
Flink CDC data ingestion jobs can synchronize schema changes from the source to the sink table, such as creating a table, adding a column, renaming a column, changing a column type, dropping a column, and dropping a table.
The sink table may not support all schema changes. You can add the schema.change.behavior option to the pipeline module to control how the sink table handles a schema change.
pipeline:
schema.change.behavior: EVOLVE
Currently, the framework supports synchronizing only the following schema change types. Schema changes that are not listed may cause job exceptions and require a stateless restart to recover.
-
Create table (
CREATE TABLE ...) -
Add column (
ALTER TABLE ... ADD COLUMN ...) -
Alter column type (
ALTER TABLE ... MODIFY COLUMN ...) -
Drop column (
ALTER TABLE ... DROP COLUMN ...) -
Rename column (
ALTER TABLE ... RENAME COLUMN ...) -
Truncate table (
TRUNCATE TABLE ...) -
Drop table (
DROP TABLE ...) -
Alter column order (
MODIFY COLUMN ... FIRST/AFTERorCHANGE COLUMN ... FIRST/AFTER)
The alter column order event can be recognized and processed in VVR 11.9 and later.
Schema evolution behavior
|
Mode |
Description |
|
LENIENT (default) |
The Flink CDC data ingestion job leniently converts schema changes and sends them to the sink table. The conversion follows these rules:
Use this mode when you want the data ingestion job to synchronize schema changes as automatically as possible. |
|
IGNORE |
The Flink CDC data ingestion job ignores all schema changes. Upstream schema changes are not applied to the sink table. Use this mode when your sink table does not support schema changes, or you do not want schema changes to occur, and you want to keep receiving data from the unchanged columns. |
|
EVOLVE |
The Flink CDC data ingestion job applies all schema changes to the sink table. If a schema change fails to be applied to the sink table, the job throws an exception and triggers a failover. Use this mode when you want the data ingestion job to synchronize schema changes as strictly as possible. Important
In this mode, if the sink cannot apply all schema change events, the job may fail over and cannot self-recover. |
|
TRY_EVOLVE |
The Flink CDC data ingestion job attempts to apply schema changes to the sink table. If the sink table cannot process a schema change, the job does not fail or restart. Instead, it attempts to handle the change by transforming the subsequent data. Use this mode when you want the data ingestion job to synchronize schema changes as strictly as possible while retaining a degree of fault tolerance. Important
In TRY_EVOLVE mode, if a schema change fails to be applied, subsequent data from the upstream may lose columns or be truncated to fit the sink table schema. |
|
EXCEPTION |
The Flink CDC data ingestion job does not allow any schema change. The job throws an exception when it receives a schema change event. Use this mode when you must ensure that no schema change occurs in the data ingestion job and only data is synchronized. |
Control schema changes at the sink
In data synchronization scenarios, Flink CDC provides fine-grained schema change management. You can use rules to control the types of schema changes that the sink receives, which prevents data loss or service interruption caused by unexpected changes.
To do this, configure the sink module with the include.schema.changes and exclude.schema.changes options.
|
Parameter |
Description |
Required |
Data type |
Default value |
Remarks |
|
include.schema.changes |
The schema changes that can be applied. |
No |
List<String> |
None |
All changes are supported by default. |
|
exclude.schema.changes |
The schema changes that cannot be applied. |
No |
List<String> |
None |
Takes precedence over |
The following table lists all configurable schema change event types:
|
Event type |
Description |
|
|
Add a column |
|
|
Change a column type |
|
|
Create a table |
|
|
Drop a column |
|
|
Drop a table |
|
|
Rename a column |
|
|
Truncate a table |
|
|
Alter column order |
Schema changes support partial matching. For example, specifying drop is equivalent to specifying both drop.column and drop.table. Specifying table is equivalent to specifying create.table, truncate.table, and drop.table.
Examples
-
Example 1: Set the schema change behavior to EVOLVE to synchronize upstream schema changes to the sink.
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: ${mysql.source.table}
server-id: 7601-7604
sink:
type: values
name: Values Sink
print.enabled: true
sink.print.logger: true
pipeline:
name: mysql to print job
schema.change.behavior: EVOLVE
-
Example 2: Apply the create-table and column events from the upstream table to the sink, and ignore drop-column events.
source:
type: mysql
name: MySQL Source
hostname: ${mysql.hostname}
port: ${mysql.port}
username: ${mysql.username}
password: ${mysql.password}
tables: ${mysql.source.table}
server-id: 7601-7604
sink:
type: values
name: Values Sink
print.enabled: true
sink.print.logger: true
include.schema.changes: [create.table, column] # Matches the CreateTable, AddColumn, AlterColumnType, RenameColumn, and DropColumn events
exclude.schema.changes: [drop.column] # Excludes the DropColumn event
pipeline:
name: mysql to print job
schema.change.behavior: EVOLVE