FEATURED · 精选文章

SeaTunnel 模式演进(Schema Evolution)实战指南:让 CDC 数据同步自动适配表结构变更

发布时间 / 2026/9/20 7:10:29
来源 / 创域科博编辑部
栏目 / 资讯中心
SeaTunnel 模式演进(Schema Evolution)实战指南:让 CDC 数据同步自动适配表结构变更 数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载导读本指南聚焦 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 源侧事件过滤还支持更细粒度的控制详见下文事件类型过滤小节。已支持的连接器源SourceMySQL-CDCOracle-CDC目标SinkJdbc-MysqlJdbc-OracleJdbc-PostgresJdbc-DamengJdbc-SqlServerStarRocksDorisPaimonElasticsearchBigQuery仅支持ADD COLUMNRedis以 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_开头也会出现相同的问题。达梦Dameng类型转换限制早期版本的达梦数据库不支持将Varchar类型字段更改为Text类型字段。启用 Schema Evolutionschema-changes.enabled在 CDC 源连接器中模式演进默认是关闭的。你需要在 CDC 连接器中配置schema-changes.enabled true该参数的定义位于 CDC 基础模块的 SourceOptions.javapublic static final OptionBoolean 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.javaschema-changes.include仅当schema-changes.enabled true时列表中列出的事件类型才会下发空列表表示所有事件类型都允许。schema-changes.exclude列出的事件类型将不会下发。该参数在include之后生效当某类型同时出现在两个列表中时exclude 优先。合法的取值由 SchemaChangeEventType.java 统一定义采用面向用户的规范名称规范名称对应事件说明add.columnSCHEMA_CHANGE_ADD_COLUMN新增列drop.columnSCHEMA_CHANGE_DROP_COLUMN删除列modify.columnSCHEMA_CHANGE_MODIFY_COLUMN修改列change.columnSCHEMA_CHANGE_CHANGE_COLUMN变更列update.columnsSCHEMA_CHANGE_UPDATE_COLUMNS列级变更的分组别名等价于所有列级变更值得注意的实现细节rename.table表重命名故意没有对外暴露。从源码注释可以看到CDC 目前对表重命名没有端到端的处理能力DDL 不会被解析成AlterTableNameEventschema 处理器将其视为空操作也没有任何 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.productsshop_b.products会写入shop_b_sink.products。如果两张源表之后都执行了ALTER TABLE products ADD COLUMN add_column1 VARCHAR(64), ADD COLUMN add_column2 INT这类 DDLSeaTunnel 会分别把 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_productsshop_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-Mysqlenv { # You can set engine configuration here parallelism 5 job.mode STREAMING checkpoint.interval 5000 read_limit.bytes_per_second7000000 read_limit.rows_per_second400 } 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-Oracleenv { # 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 - StarRocksenv { # 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_modeRECREATE_SCHEMA data_save_modeDROP_DATA save_mode_create_template CREATE TABLE IF NOT EXISTS shop.${table_name} ( ${rowtype_primary_key}, ${rowtype_fields} ) ENGINEOLAP 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 - Dorisenv { # 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-Postgresenv { # You can set engine configuration here parallelism 5 job.mode STREAMING checkpoint.interval 5000 read_limit.bytes_per_second7000000 read_limit.rows_per_second400 } 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_second7000000 read_limit.rows_per_second400 } 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-SqlServerenv { # You can set engine configuration here parallelism 5 job.mode STREAMING checkpoint.interval 5000 read_limit.bytes_per_second7000000 read_limit.rows_per_second400 } 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大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载相关推荐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),仅供参考
RELATED — 相关阅读

相关资讯

LATEST — 最新资讯

最新发布

TODAY — 本日精选

新闻

WEEKLY — 本周精选

新闻

MONTHLY — 本月精选

新闻