- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
导读
本指南聚焦 Apache SeaTunnel 的模式演进(Schema Evolution)能力:当上游数据表执行ALTER TABLE等 DDL 变更后,数据同步任务无需停机、无需手工改配置,即可自动感知新表结构并继续同步。文章基于docs/zh/introduction/configuration/schema-evolution.md展开,并结合仓库源码说明schema-changes.enabled等参数的底层实现与事件过滤机制。读完本文,你将掌握:模式演进的适用引擎与连接器范围、已知限制、多库多表下的路由配置,以及七组可直接复用的 CDC 到 JDBC / StarRocks / Doris / Paimon 等目标端的完整 HOCON 配置。
什么是模式演进(Schema Evolution)
模式演进是指数据表的 Schema(表结构)可以动态改变,数据同步任务能够自动适应新的表结构变化,而无需任何人工干预。在传统数据集成任务中,一旦上游表ADD COLUMN或修改了字段类型,同步作业往往因写入字段与下游表结构不匹配而失败,需要运维人员手动调整任务并重启。开启模式演进后,SeaTunnel 会拦截上游的 DDL 事件并将其转换为 schema 变更事件下发给下游,由下游连接器自动执行对应的 DDL,实现表结构的"零停机跟随"。
从源码看,这一机制建立在 SeaTunnel 的事件体系之上:CDC 源将 DDL 解析为EventType中的SCHEMA_CHANGE_ADD_COLUMN、SCHEMA_CHANGE_DROP_COLUMN等事件(见 EventType 相关枚举与 SchemaChangeEventType 映射),经引擎路由到支持该变更类型的 Sink 连接器执行。
支持范围一览
已支持的引擎
| 引擎 | 支持状态 |
|---|---|
| SeaTunnel Zeta | ✅ 支持 |
已支持的模式变更事件类型
ADD COLUMN:新增列DROP COLUMN:删除列RENAME COLUMN:重命名列MODIFY COLUMN:修改列(类型、约束等)
需要说明的是,上述列表是连接器层面的能力范围。在 CDC 源侧,事件过滤还支持更细粒度的控制(详见下文"事件类型过滤"小节)。
已支持的连接器
源(Source)
- MySQL-CDC
- Oracle-CDC
目标(Sink)
- Jdbc-Mysql
- Jdbc-Oracle
- Jdbc-Postgres
- Jdbc-Dameng
- Jdbc-SqlServer
- StarRocks
- Doris
- Paimon
- Elasticsearch
- BigQuery(仅支持
ADD COLUMN) - Redis
以 Doris 为例,其 Sink 在supports()方法中声明支持ADD_COLUMN、DROP_COLUMN、RENAME_COLUMN、UPDATE_COLUMN四类变更(见 DorisSink.java)。每个连接器通过实现自身的 schema 变更处理器来决定如何处理具体 DDL,这也解释了为什么不同连接器支持的事件类型存在差异。
已知限制与注意事项
使用模式演进前,请务必了解以下边界:
- 不支持与 Transform 叠加使用:目前模式演进不支持 transform。若任务中配置了 transform,则无法同时启用 schema 演进。
- 跨库类型的列默认值缺失:不同类型数据库之间的模式演进(如
Oracle-CDC -> Jdbc-Mysql)目前不支持 DDL 中列的默认值。 - Oracle-CDC 的账号与表名约束:
- 使用 Oracle-CDC 时,不能使用用户名
SYS或SYSTEM修改表结构,否则 DDL 事件会被过滤,导致模式演进不起作用; - 如果表名以
ORA_TEMP_开头,也会出现相同的问题。
- 使用 Oracle-CDC 时,不能使用用户名
- 达梦(Dameng)类型转换限制:早期版本的达梦数据库不支持将
Varchar类型字段更改为Text类型字段。
启用 Schema Evolution:schema-changes.enabled
在 CDC 源连接器中,模式演进默认是关闭的。你需要在 CDC 连接器中配置:
schema-changes.enabled = true该参数的定义位于 CDC 基础模块的 SourceOptions.java:
public static final Option<Boolean> SCHEMA_CHANGES_ENABLED = Options.key("schema-changes.enabled") .booleanType() .defaultValue(false) .withDescription( "Enable send schema change events, by default is false. If set to true, the schema changes will be sent to downstream.");即:schema-changes.enabled默认值为false,只有在显式设置为true后,DDL 事件才会被转换为 schema 变更事件发送给下游 Sink 执行。
进阶:事件类型过滤(schema-changes.include / exclude)
除了总开关,源码中还提供了两个细粒度过滤参数,可以精确控制哪些类型的 DDL 事件允许下发(见 SourceOptions.java):
schema-changes.include:仅当schema-changes.enabled = true时,列表中列出的事件类型才会下发;空列表表示所有事件类型都允许。schema-changes.exclude:列出的事件类型将不会下发。该参数在include之后生效;当某类型同时出现在两个列表中时,exclude 优先。
合法的取值由 SchemaChangeEventType.java 统一定义,采用面向用户的规范名称:
| 规范名称 | 对应事件 | 说明 |
|---|---|---|
add.column | SCHEMA_CHANGE_ADD_COLUMN | 新增列 |
drop.column | SCHEMA_CHANGE_DROP_COLUMN | 删除列 |
modify.column | SCHEMA_CHANGE_MODIFY_COLUMN | 修改列 |
change.column | SCHEMA_CHANGE_CHANGE_COLUMN | 变更列 |
update.columns | SCHEMA_CHANGE_UPDATE_COLUMNS | 列级变更的分组别名,等价于所有列级变更 |
值得注意的实现细节:rename.table(表重命名)故意没有对外暴露。从源码注释可以看到,CDC 目前对表重命名没有端到端的处理能力:DDL 不会被解析成AlterTableNameEvent,schema 处理器将其视为空操作,也没有任何 Sink 会执行它。如果将其暴露为可过滤名称,等于宣传一个并不存在的能力(见 SchemaChangeEventType.java 的注释)。
示例:只允许新增列、拒绝删除列:
source { MySQL-CDC { ... schema-changes.enabled = true schema-changes.include = ["add.column"] schema-changes.exclude = ["drop.column"] } }多库多表路由下的模式演进
只要每张上游表都能稳定映射到一个明确的物理下游表,模式演进就可以和多库多表任务一起工作。SeaTunnel 会在连接器启动前完成 Sink 占位符替换,因此你可以结合 Sink 参数占位符 中的${database_name}、${schema_name}、${table_name}做路由。
推荐做法
- 如果希望不同上游库的表彼此隔离,请把它们路由到不同的物理下游表。
- 如果需要并行写入,可继续开启
multi_table_sink_replica;模式变更会按最终渲染出的物理下游表维度协调执行。 - 如果你有意把多张上游表写入同一张物理下游表,请自行保证这些表的 schema 兼容,并确保主键不会冲突。
示例一:不同源库中的同名表 -> 不同下游库中的同名表
source { MySQL-CDC { database-names = ["shop_a", "shop_b"] table-names = ["shop_a.products", "shop_b.products"] url = "jdbc:mysql://mysql-host:3306" schema-changes.enabled = true } } sink { jdbc { url = "jdbc:mysql://mysql-host:3306" driver = "com.mysql.cj.jdbc.Driver" user = "root" password = "123456" generate_sink_sql = true database = "${database_name}_sink" table = "${table_name}" primary_keys = ["id"] multi_table_sink_replica = 2 } }在这个例子里,shop_a.products会写入shop_a_sink.products,shop_b.products会写入shop_b_sink.products。
如果两张源表之后都执行了ALTER TABLE products ADD COLUMN add_column1 VARCHAR(64), ADD COLUMN add_column2 INT这类 DDL,SeaTunnel 会分别把 schema 变更应用到shop_a_sink.products和shop_b_sink.products,并继续保证每张下游表只接收自己所属源库的数据。
示例二:写入同一个下游库,但拆成不同下游表
sink { jdbc { url = "jdbc:mysql://mysql-host:3306" driver = "com.mysql.cj.jdbc.Driver" user = "root" password = "123456" generate_sink_sql = true database = "ods" table = "${database_name}_${table_name}" primary_keys = ["id"] } }在这个例子里,shop_a.products会写入ods.shop_a_products,shop_b.products会写入ods.shop_b_products。
示例三:用通配符捕获多库多表
source { MySQL-CDC { table-pattern = "sales_.*\\..*" url = "jdbc:mysql://mysql-host:3306" schema-changes.enabled = true } } sink { jdbc { url = "jdbc:mysql://mysql-host:3306" driver = "com.mysql.cj.jdbc.Driver" user = "root" password = "123456" generate_sink_sql = true database = "ods" table = "${database_name}_${table_name}" primary_keys = ["${primary_key}"] } }关于占位符与主键展开
${primary_key}占位符会在连接器启动前被展开为上游表元数据中定义的所有主键列。例如上游主键为(f1, f2)时,primary_keys = ["${primary_key}"]会被展开为primary_keys = ["f1", "f2"]。当前不支持将${primary_key}与静态列名混合使用(如primary_keys = ["${primary_key}", "tenant_id"]),只有当${primary_key}是列表中的唯一元素时才会执行列表占位符替换(详见 Sink 参数占位符)。另外注意:mysql源不包含${schema_name}元数据、oracle源不包含${database_name}元数据,若占位符未替换需检查上游元数据是否提供了对应信息。
完整配置示例
以下示例均为仓库中经过端到端验证的真实配置,可直接作为任务模板参考。
Mysql-CDC -> Jdbc-Mysql
env { # You can set engine configuration here parallelism = 5 job.mode = "STREAMING" checkpoint.interval = 5000 read_limit.bytes_per_second=7000000 read_limit.rows_per_second=400 } source { MySQL-CDC { server-id = 5652-5657 username = "st_user_source" password = "mysqlpw" table-names = ["shop.products"] url = "jdbc:mysql://mysql_cdc_e2e:3306/shop" schema-changes.enabled = true } } sink { jdbc { url = "jdbc:mysql://mysql_cdc_e2e:3306/shop" driver = "com.mysql.cj.jdbc.Driver" user = "st_user_sink" password = "mysqlpw" generate_sink_sql = true database = shop table = mysql_cdc_e2e_sink_table_with_schema_change_exactly_once primary_keys = ["id"] is_exactly_once = true xa_data_source_class_name = "com.mysql.cj.jdbc.MysqlXADataSource" } }要点说明:该示例开启了is_exactly_once = true,通过xa_data_source_class_name指定 MySQL XA 数据源实现精确一次语义,与 schema 演进配合使用。generate_sink_sql = true表示由 Sink 根据上游元数据自动生成建表/写入 SQL,这是模式演进能够自动适配新列的前提。
Oracle-CDC -> Jdbc-Oracle
env { # You can set engine configuration here parallelism = 1 job.mode = "STREAMING" checkpoint.interval = 5000 } source { # This is a example source plugin **only for test and demonstrate the feature source plugin** Oracle-CDC { plugin_output = "customers" username = "dbzuser" password = "dbz" database-names = ["ORCLCDB"] schema-names = ["DEBEZIUM"] table-names = ["ORCLCDB.DEBEZIUM.FULL_TYPES"] url = "jdbc:oracle:thin:@oracle-host:1521/ORCLCDB" source.reader.close.timeout = 120000 connection.pool.size = 1 schema-changes.enabled = true } } sink { Jdbc { plugin_input = "customers" driver = "oracle.jdbc.driver.OracleDriver" url = "jdbc:oracle:thin:@oracle-host:1521/ORCLCDB" user = "dbzuser" password = "dbz" generate_sink_sql = true database = "ORCLCDB" table = "DEBEZIUM.FULL_TYPES_SINK" batch_size = 1 primary_keys = ["ID"] connection.pool.size = 1 } }要点说明:plugin_output与plugin_input成对使用,用于将 CDC 源的数据流显式路由到指定的 Sink 输入。Oracle 侧需要注意文档前面提到的限制:不要使用SYS/SYSTEM账号执行 DDL,且表名不能以ORA_TEMP_开头。
Oracle-CDC -> Jdbc-Mysql(跨数据库类型)
env { # You can set engine configuration here parallelism = 1 job.mode = "STREAMING" checkpoint.interval = 5000 } source { # This is a example source plugin **only for test and demonstrate the feature source plugin** Oracle-CDC { plugin_output = "customers" username = "dbzuser" password = "dbz" database-names = ["ORCLCDB"] schema-names = ["DEBEZIUM"] table-names = ["ORCLCDB.DEBEZIUM.FULL_TYPES"] url = "jdbc:oracle:thin:@oracle-host:1521/ORCLCDB" source.reader.close.timeout = 120000 connection.pool.size = 1 schema-changes.enabled = true } } sink { jdbc { plugin_input = "customers" url = "jdbc:mysql://oracle-host:3306/oracle_sink" driver = "com.mysql.cj.jdbc.Driver" user = "st_user_sink" password = "mysqlpw" generate_sink_sql = true # You need to configure both database and table database = oracle_sink table = oracle_cdc_2_mysql_sink_table primary_keys = ["ID"] } }跨库类型场景下需要同时显式配置database和table;同时注意此类链路中 DDL 的列默认值暂不支持(见前文限制说明)。
MySQL-CDC -> StarRocks
env { # You can set engine configuration here parallelism = 1 job.mode = "STREAMING" checkpoint.interval = 5000 } source { MySQL-CDC { username = "st_user_source" password = "mysqlpw" table-names = ["shop.products"] url = "jdbc:mysql://mysql_cdc_e2e:3306/shop" schema-changes.enabled = true } } sink { StarRocks { nodeUrls = ["starrocks_cdc_e2e:8030"] username = "root" password = "" database = "shop" table = "${table_name}" url = "jdbc:mysql://starrocks_cdc_e2e:9030/shop" max_retries = 3 enable_upsert_delete = true schema_save_mode="RECREATE_SCHEMA" data_save_mode="DROP_DATA" save_mode_create_template = """ CREATE TABLE IF NOT EXISTS shop.`${table_name}` ( ${rowtype_primary_key}, ${rowtype_fields} ) ENGINE=OLAP PRIMARY KEY (${rowtype_primary_key}) DISTRIBUTED BY HASH (${rowtype_primary_key}) PROPERTIES ( "replication_num" = "1", "in_memory" = "false", "enable_persistent_index" = "true", "replicated_storage" = "true", "compression" = "LZ4" ) """ } }要点说明:StarRocks 场景下通过${table_name}占位符动态路由目标表;save_mode_create_template中${rowtype_primary_key}与${rowtype_fields}是建表模板的内置占位符,分别渲染为主键字段与全量字段列表,schema 变更后新建的表会按最新字段集合生成。
MySQL-CDC -> Doris
env { # You can set engine configuration here parallelism = 1 job.mode = "STREAMING" checkpoint.interval = 5000 } source { MySQL-CDC { server-id = 5652-5657 username = "st_user_source" password = "mysqlpw" table-names = ["shop.products"] url = "jdbc:mysql://mysql_cdc_e2e:3306/shop" schema-changes.enabled = true } } sink { Doris { fenodes = "doris_e2e:8030" username = "root" password = "" database = "shop" table = "products" sink.label-prefix = "test-cdc" sink.enable-2pc = "true" sink.enable-delete = "true" doris.config { format = "json" read_json_by_line = "true" } } }注意(schema 演进 + 2PC):当
sink.enable-2pc = "true"时,Doris schema 演进仅支持format = "json",因为 JSON load 会按列名匹配。CSV 等位置敏感格式在启用 2PC 的 schema 演进场景下会被运行时拒绝。请使用format = "json",或设置sink.enable-2pc = "false",让 sink 可以在应用 DDL 前先 flush 已缓冲的数据。
Doris 侧的模式演进由SchemaChangeManager统一调度(见 connector-doris 的 schema 目录),DorisSink 声明支持ADD_COLUMN、DROP_COLUMN、RENAME_COLUMN、UPDATE_COLUMN四类变更。2PC 场景下使用 JSON 格式是因为按列名匹配才能与 DDL 后的新 schema 保持一致。
MySQL-CDC -> Jdbc-Postgres
env { # You can set engine configuration here parallelism = 5 job.mode = "STREAMING" checkpoint.interval = 5000 read_limit.bytes_per_second=7000000 read_limit.rows_per_second=400 } source { MySQL-CDC { server-id = 5652-5657 username = "st_user_source" password = "mysqlpw" table-names = ["shop.products"] url = "jdbc:mysql://mysql_cdc_e2e:3306/shop" schema-changes.enabled = true } } sink { jdbc { url = "jdbc:postgresql://postgresql:5432/shop" driver = "org.postgresql.Driver" user = "postgres" password = "postgres" generate_sink_sql = true database = shop table = "public.sink_table_with_schema_change" primary_keys = ["id"] # Validate ddl update for sink writer multi replica multi_table_sink_replica = 2 } }MySQL-CDC -> Jdbc-Dameng(达梦)
env { # You can set engine configuration here parallelism = 5 job.mode = "STREAMING" checkpoint.interval = 5000 read_limit.bytes_per_second=7000000 read_limit.rows_per_second=400 } source { MySQL-CDC { server-id = 5652-5657 username = "st_user_source" password = "mysqlpw" table-names = ["shop.products"] url = "jdbc:mysql://mysql_cdc_e2e:3306/shop" schema-changes.enabled = true } } sink { jdbc { url = "jdbc:dm://e2e_dmdb:5236" driver = "dm.jdbc.driver.DmDriver" connection_check_timeout_sec = 1000 user = "SYSDBA" password = "SYSDBA" generate_sink_sql = true database = "DAMENG" table = "SYSDBA.sink_table_with_schema_change" primary_keys = ["id"] # Validate ddl update for sink writer multi replica multi_table_sink_replica = 2 } }MySQL-CDC -> Jdbc-SqlServer
env { # You can set engine configuration here parallelism = 5 job.mode = "STREAMING" checkpoint.interval = 5000 read_limit.bytes_per_second=7000000 read_limit.rows_per_second=400 } source { MySQL-CDC { server-id = 5652-5657 username = "st_user_source" password = "mysqlpw" table-names = ["shop.products"] url = "jdbc:mysql://mysql_cdc_e2e:3306/shop" schema-changes.enabled = true } } sink { jdbc { url = "jdbc:sqlserver://e2e_sqlserver:1433" driver = "com.microsoft.sqlserver.jdbc.SQLServerDriver" user = "sa" password = "paanssy1234$" generate_sink_sql = true database = master table = "dbo.sink_table_with_schema_change" primary_keys = ["id"] # Validate ddl update for sink writer multi replica multi_table_sink_replica = 2 } }源码级原理:schema 变更如何被协调执行
从实现角度看,模式演进并不仅仅是"把 DDL 透传下去",还涉及引擎与 Sink 的多层协作:
- 源侧解析与过滤:CDC 源捕获 DDL 后,依据
schema-changes.enabled、schema-changes.include、schema-changes.exclude决定是否下发、下发哪些事件类型(SourceOptions.java)。 - 事件规范化:DDL 被解析为 SeaTunnel 统一的 schema 变更事件对象,内部通过
EventType枚举区分ADD_COLUMN/DROP_COLUMN/MODIFY_COLUMN/CHANGE_COLUMN/UPDATE_COLUMNS等类型(SchemaChangeEventType.java)。 - Sink 声明能力并执行:每个支持模式演进的 Sink 通过
supports()声明自己能处理的变更类型(如 DorisSink.java),随后在收到事件时转换为对应数据库的 ALTER 语句并执行。 - 多副本协调:开启
multi_table_sink_replica后,schema 变更会按最终渲染出的物理下游表维度协调执行,确保同一物理表的多个 writer 副本都完成 DDL 后才继续写入。
仓库中的回归测试验证了这一协调逻辑。例如 JdbcSinkWriterSharedPhysicalTableSchemaChangeTest.java 使用真实 SQLite 数据库验证了"两张上游表折叠到同一物理下游表"的场景:上游表 A 执行 DROP COLUMN 后,协调器会先重建共享物理表对应的全部 writer,再让上游表 B 按新 schema 继续写入。这正好对应了文档中"把多张上游表写入同一张物理下游表时需自行保证 schema 兼容"的建议。
端到端层面,仓库在seatunnel-e2e中提供了大量可运行验证,例如 mysqlcdc_to_mysql_with_schema_change.conf、mysqlcdc_to_mysql_with_multi_db_same_name_schema_change.conf(多库同名表路由)等,以及对应的MysqlCDCWithSchemaChangeIT集成测试类,可作为模式演进配置的权威参考。
小结与建议
SeaTunnel 的模式演进能力把"表结构变更"从运维事故变成了可自动处理的事件流。落地使用时请记住几条核心原则:
- 总开关默认关闭:必须在 CDC 源侧显式设置
schema-changes.enabled = true;如需更细粒度控制,配合schema-changes.include/schema-changes.exclude使用。 - 确认连接器能力:不同目标端支持的事件类型不同(如 BigQuery 仅支持
ADD COLUMN),且表重命名(rename.table)目前端到端不可用。 - 规划好路由:多库多表场景下优先让每张上游表稳定映射到独立物理下游表,必要时开启
multi_table_sink_replica提升写入并行度;共享同一物理表的任务需自行保证 schema 兼容与主键不冲突。 - 规避已知雷区:不要与 transform 混用;Oracle-CDC 避免使用
SYS/SYSTEM账号及ORA_TEMP_前缀表;Doris 2PC 场景必须使用 JSON 格式;跨数据库类型(如 Oracle -> MySQL)的 DDL 列默认值暂不支持。
结合仓库中的端到端配置与回归测试,你可以直接复用上文示例,快速搭建一套能"随表而变"的实时数据同步管道。
- 数据集成
- ETL
- 大数据
- 批处理
- 流处理
- 变更数据捕获
【免费下载链接】seatunnel
SeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.
相关推荐
Flink CDC Schema演化:5种模式实现数据库表结构变更自动同步
Flink CDC Schema演化:5种模式实现数据库表结构变更自动同步 Flink CDC Schema演化是流式数据集成中的关键技术,能够自动同步上游数据
后端数据集成大数据流处理变更数据捕获数据同步SeaTunnel Schema Evolution 实战指南:基于 CDC 的库表结构自动演进与多表路由
SeaTunnel Schema Evolution 实战指南:基于 CDC 的库表结构自动演进与多表路由 Schema Evolution(表结构演进)是 S
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Paimon Sink 连接器实战指南:CDC 同步、自动建表、Schema Evolution 与多表写入全解析
SeaTunnel Paimon Sink 连接器实战指南:CDC 同步、自动建表、Schema Evolution 与多表写入全解析 导读 本文以 docs/
数据集成ETL大数据批处理流处理变更数据捕获
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考