SeaTunnel MySQL CDC 到 Doris 实时同步实战:从全量快照到增量变更的端到端配置指南
2026/9/20 2:27:03 网站建设 项目流程

SeaTunnel MySQL CDC 到 Doris 实时同步实战:从全量快照到增量变更的端到端配置指南

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

导读

本文是一份面向 SeaTunnel 开发者的端到端实战指南,完整讲解如何用 SeaTunnel 把 MySQL 的行级变更(CDC)持续同步到 Apache Doris,让 Doris 中的数据始终与 MySQL 保持最新。读完本文,你将掌握这条链路的插件与驱动安装、MySQL binlog 与 CDC 账号准备、最小 HOCON 配置编写、流式作业的启动与验证,以及server-id、Stream Load 标签、删除事件传播等高频坑位的排查方法,并理解背后的 MySQL CDC Source 与 Doris Sink 源码级实现原理。


链路概述:一条 MySQL 变更事件持续流向 Doris 的通道

当业务库发生 INSERT / UPDATE / DELETE 时,MySQL 会把这些行级变更写入 binlog。SeaTunnel 的MySQL-CDCSource 连接器以"全量快照 + 增量 binlog"两阶段方式读取这些变更(startup.mode = "initial"时先做一致性快照,快照完成后自动切换到增量读取,切换窗口内不丢事件,详见 MySQL CDC 文档),随后由DorisSink 连接器通过 Doris 的 Stream Load 批量导入机制写入目标表。由于 Doris Sink 内部以 Stream Load 批量缓存和导入为实现基础(见 Doris Sink 文档),配合sink.enable-delete即可把 UPDATE / DELETE 事件正确回放到 Doris 的 Unique Key 模型中。

整条链路由以下三部分组成:

组成连接器/模块职责
数据源MySQL-CDC(文档)捕获 MySQL 快照与 binlog 增量变更
引擎SeaTunnel Zeta(也可用 Flink / Spark)编排流式任务、checkpoint 与状态管理
目标端Doris(文档)通过 Stream Load 将变更写入 Doris

一、前置条件:四步准备让链路具备运行基础

在写配置之前,需要先完成环境、插件、源库与目标库的准备工作。官方教程要求先完成 跑第一个任务,确认本地基础链路正常,再进入本教程。

1.1 安装连接器插件并放入 JDBC 驱动

从 2.2.0-beta 版本开始,SeaTunnel 二进制发行包不再默认携带所有连接器依赖,需要按需安装。按 部署 > 下载连接器插件 的说明,把config/plugin_config修改为只保留这条链路需要的两个插件:

--seatunnel-connectors-- connector-cdc-mysql connector-doris --end--

注意:仓库中默认的 config/plugin_config 使用--connectors-v2--区块列出了全部插件(其中已经包含connector-cdc-mysqlconnector-doris),而教程示例使用的是--seatunnel-connectors--区块写法。无论哪种区块,核心原则都是"只安装任务实际需要的插件",以减小体积、加快启动。

随后执行安装脚本并确认插件落盘:

cd "${SEATUNNEL_HOME}" sh bin/install-plugin.sh ls connectors | rg 'connector-(cdc-mysql|doris)'

接着处理 MySQL JDBC 驱动。由于 MySQL CDC 需要建立 JDBC 连接执行快照与元数据查询,必须把mysql-connector-java驱动放到引擎实际加载的位置:

  • SeaTunnel Zeta:放进${SEATUNNEL_HOME}/lib,并确认 jar 已落盘:
ls "${SEATUNNEL_HOME}/lib" | rg 'mysql-connector'
  • Flink / Spark:把同一个驱动 jar 放到对应引擎实际加载的插件目录(如${SEATUNNEL_HOME}/plugins/)。

1.2 准备 MySQL 源表:稳定主键是正确回放变更的前提

教程的验证依赖稳定主键——有了主键,Doris 自动建表时才能生成 Unique Key,下游才能正确区分并回放 UPDATE 与 DELETE 事件。执行以下 SQL:

CREATE DATABASE IF NOT EXISTS inventory; CREATE TABLE IF NOT EXISTS inventory.orders ( id BIGINT PRIMARY KEY, order_status VARCHAR(32), amount DECIMAL(10, 2), updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP ); INSERT INTO inventory.orders (id, order_status, amount, updated_at) VALUES (1001, 'CREATED', 19.99, NOW()), (1002, 'CREATED', 29.99, NOW());

1.3 创建 CDC 用户并授权

MySQL CDC 连接器底层依赖 Debezium 嵌入式引擎,需要用户具备读取 binlog 与执行快照的权限。按 MySQL CDC Source 文档 创建用户并授权:

CREATE USER IF NOT EXISTS 'st_user_source'@'%' IDENTIFIED BY 'mysqlpw'; GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'st_user_source'@'%'; FLUSH PRIVILEGES;

各权限的用途可以结合源码语义理解:SELECT用于快照阶段读取全量数据,RELOAD用于一致性快照的锁表/解锁,SHOW DATABASES用于枚举库表,REPLICATION SLAVEREPLICATION CLIENT用于订阅和查询 binlog 位点。

1.4 检查并开启 MySQL binlog

MySQL CDC 完全依赖 binlog 工作,先检查三个关键变量:

SHOW VARIABLES WHERE variable_name IN ('log_bin', 'binlog_format', 'binlog_row_image');

期望值是log_bin = ONbinlog_format = ROWbinlog_row_image = FULL。如果未配置,修改my.cnf并重启 MySQL:

[mysqld] server-id = 223344 log_bin = mysql-bin binlog_format = ROW binlog_row_image = FULL

扩展提示:MySQL 5.6+ 还建议开启 GTID(gtid_mode = onenforce_gtid_consistency = on),这并非强制要求,但能让故障恢复与位点定位更可靠;MySQL 8.0+ 默认 GTID 为 OFF 也可正常工作(详见 MySQL CDC 文档)。另外,对大型数据库做初始快照时,若连接可能超时,可在 MySQL 配置中调大interactive_timeoutwait_timeout

1.5 准备 Doris 目标库

本教程保留schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST",因此第一次启动时会利用 MySQL 主键信息在 Doris 自动创建sync_demo.orders,你只需先建好数据库:

CREATE DATABASE IF NOT EXISTS sync_demo;

二、最小配置详解:每个参数背后的实现语义

将下面的配置保存为config/mysql-cdc-to-doris.conf。这是本教程的核心配置,下面逐块拆解:

env { parallelism = 1 job.mode = "STREAMING" checkpoint.interval = 5000 } source { MySQL-CDC { plugin_output = "orders_cdc" parallelism = 1 startup.mode = "initial" server-id = 5652 username = "st_user_source" password = "mysqlpw" table-names = ["inventory.orders"] url = "jdbc:mysql://mysql:3306/inventory" } } sink { Doris { plugin_input = "orders_cdc" fenodes = "doris-fe:8030" username = "root" password = "" database = "sync_demo" table = "orders" sink.label-prefix = "orders-cdc" sink.enable-delete = true schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST" doris.config = { format = "csv" column_separator = "," } } }

2.1 env 块:流式作业的关键开关

参数说明
parallelism1单并发。CDC 场景下若要多并发读多张表,需配合足够大的server-id范围
job.mode"STREAMING"声明为流式任务,作业将持续运行并处理增量事件
checkpoint.interval5000每 5 秒触发一次 checkpoint,决定故障恢复粒度与数据可见性节奏

2.2 source 块:MySQL CDC 核心参数

教程用到的参数语义如下(完整参数表见 MySQL CDC 配置参数选项):

  • plugin_output = "orders_cdc":给下游 sink 使用的数据流名称,必须与 Doris sink 的plugin_input一致。
  • startup.mode = "initial":启动时先同步历史全量数据,再无缝切换同步增量数据。这是 CDC 最常见的"先全量后增量"模式。其他可选值包括earliest(最早偏移量)、latest(跳过快照、只读新变更)、specific(从指定 binlog 文件+位置启动)、timestamp(从指定时间戳启动)。
  • server-id = 5652:此 CDC 读取器使用的数字 ID。每个 ID 在 MySQL 集群中必须唯一——若与其他副本或 CDC 作业冲突,MySQL 会直接断开其中一个客户端连接。支持配置范围如5652-5657以支撑多并发/多表读取;未配置时 SeaTunnel 会随机生成,但生产环境建议显式指定。
  • table-names = ["inventory.orders"]:要监控的表,表名必须带库名前缀。也可改用table-pattern正则批量匹配,但两者只能二选一。
  • url:JDBC 连接串,示例为jdbc:mysql://mysql:3306/inventory(容器环境内用服务名mysql,本机则替换为localhost或 IP)。

其他值得了解的常用参数:snapshot.split.size(默认 8096,控制快照分块行数)、snapshot.fetch.size(默认 1024,单次抓取行数)、server-time-zone(默认 UTC,建议与 MySQL 服务器一致)、exactly_once(默认 false,启用后配合下游实现精确一次)、debezium(透传 Debezium 属性,如心跳配置)。

2.3 sink 块:Doris Stream Load 核心参数

  • plugin_input = "orders_cdc":与 source 的plugin_output对接。
  • fenodes = "doris-fe:8030":Doris FE 地址,格式为"fe_ip:fe_http_port,..."。可配置多个 FE 以逗号分隔。
  • database/table:目标库表。多表 CDC 场景还支持${database_name}/${table_name}占位符动态映射上游库表名。
  • sink.label-prefix = "orders-cdc":Stream Load 的标签前缀,用于保证导入幂等与 2PC 语义。多个运行中的任务若复用同一个 label-prefix,会导致 Doris Stream Load 标签冲突。
  • sink.enable-delete = true:开启删除事件传播,把 CDC 的 DELETE 事件回放到 Doris。要求 Doris 目标表开启批量删除且使用Unique Key模型。
  • schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST":表不存在时自动建表,存在则跳过。Doris Sink 支持四种模式:RECREATE_SCHEMA(存在则删表重建)、CREATE_SCHEMA_WHEN_NOT_EXIST(默认值,不存在才创建)、ERROR_WHEN_SCHEMA_NOT_EXIST(不存在即报错)、IGNORE(忽略建表处理)。
  • doris.config:透传给 Doris Stream Load 的数据描述参数。本配置使用 CSV 格式与逗号分隔符;也常用format = "json"+read_json_by_line = "true"的 JSON 逐行导入方式。

扩展:若同时开启sink.enable-2pc = true,Doris Sink 将通过 Stream Load 两阶段提交实现精确一次(exactly-once),但此时sink.buffer-size将失去作用,只有 checkpoint 才能触发提交(详见 Doris Sink 文档)。sink.label-prefix在 2PC 场景下必须全局唯一,否则重试或重启时会出现 "Label already exists" 错误;可在 label 中加入时间戳/唯一标识,或用CANCEL LOAD WHERE LABEL LIKE 'your-prefix%'清理残留事务。

2.4 自动建表背后的模板机制

schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"的自动建表由save_mode_create_template驱动,连接器会根据上游类型与 schema 自动填充模板占位符。默认模板如下(详见 Doris Sink 文档):

CREATE TABLE IF NOT EXISTS `${database}`.`${table_name}` ( ${rowtype_primary_key}, ${rowtype_fields} ) ENGINE=OLAP UNIQUE KEY (${rowtype_primary_key}) COMMENT '${comment}' DISTRIBUTED BY HASH (${rowtype_primary_key}) PROPERTIES ( "replication_allocation" = "tag.location.default: 1", "in_memory" = "false", "storage_format" = "V2", "disable_auto_compaction" = "false" )

注意默认模板使用UNIQUE KEY模型并保留主键占位符——这正是 CDC 的 UPDATE / DELETE 事件能在 Doris 正确回放的结构基础。占位符含义:${database}/${table_name}取上游库表名,${rowtype_fields}展开全部字段,${rowtype_primary_key}展开主键列表(可能为多列),${comment}取表注释。可在模板中自定义字段(如添加id字段并设置"replication_num" = "1"),连接器会自动从上游填充其余字段。


三、运行任务与端到端验证

3.1 启动流式作业

配置保存为config/mysql-cdc-to-doris.conf后,用本地模式启动 SeaTunnel:

cd "${SEATUNNEL_HOME}" ./bin/seatunnel.sh --config ./config/mysql-cdc-to-doris.conf -m local

这是一条流式 CDC 作业,任务启动后会持续运行——执行下面的验证 SQL 时,任务必须保持在运行状态。

3.2 执行变更并核对结果

  1. 启动作业,等待首轮全量快照完成(此时10011002两行历史数据已同步到 Doris)。
  2. 在 MySQL 中依次执行插入、更新、删除:
INSERT INTO inventory.orders (id, order_status, amount, updated_at) VALUES (1003, 'CREATED', 39.99, NOW()); UPDATE inventory.orders SET order_status = 'PAID', updated_at = NOW() WHERE id = 1001; DELETE FROM inventory.orders WHERE id = 1002;
  1. 在 Doris 中查询最新结果:
SELECT COUNT(*) FROM sync_demo.orders; SELECT id, order_status, amount FROM sync_demo.orders ORDER BY id;

期望结果:1001的状态变为PAID1003成功插入,1002已被删除。若新增、更新、删除三类事件都能体现在 Doris 中,即证明这条 CDC 链路已打通。

验证过程的实现依据:MySQL CDC 会为每行变更携带行级操作类型(INSERT / UPDATE_BEFORE / UPDATE_AFTER / DELETE),Doris Sink 的sink.enable-delete负责把 DELETE 事件转换为 Doris 的批量删除语义。仓库的 E2E 测试 write-cdc-changelog-to-doris.conf 与 mysqlcdc_to_doris_with_schema_change.conf 使用的正是与本教程几乎相同的MySQL-CDC + Doris组合(含sink.enable-2pcsink.enable-deleteformat = "csv"/"json"等),对应的测试类DorisCDCSinkITDorisSchemaChangeIT会构造完整的增删改事件并断言 Doris 侧结果,可作为端到端行为的权威参考。另外 MysqlCDCIT.java 也验证了 MySQL CDC 的快照 + 增量读取能力。


四、常见坑位与排查清单

教程在"常见坑"一节列出的问题,结合源码文档可以整理为以下排查清单:

  1. MySQL 没开 binlog,或 binlog 不是ROW格式。CDC 依赖行级 binlog,务必确认log_bin = ONbinlog_format = ROWbinlog_row_image = FULL
  2. CDC 用户缺少复制相关权限。缺少REPLICATION SLAVE/REPLICATION CLIENT会导致无法订阅 binlog,需重新授权并FLUSH PRIVILEGES
  3. server-id与 MySQL 其他副本或 CDC 作业冲突。每个读取器 ID 在集群中必须唯一,冲突会导致 MySQL 断开连接。多并发/多表任务请配置足够大的 ID 范围,例如一个任务5400-5600、另一个5601-5800,互不重叠。
  4. 多个运行中的任务复用了同一个sink.label-prefix,导致 Doris Stream Load 标签冲突(典型报错 "Label already exists")。为每个任务设置唯一前缀,或在 2PC 前缀中加入时间戳/唯一标识。
  5. 开启了删除同步,但 Doris 目标表模型不支持预期的删除行为sink.enable-delete仅支持 Doris 的Unique Key模型(0.15+ 默认开启批量删除),若建表模板被改成了其他模型,DELETE 事件将无法正确回放。
  6. 源表没有稳定主键,导致自动建表和下游 upsert 结果不稳定。Doris 自动建表依赖上游主键生成 Unique Key;无主键表可退而求其次,通过table-names-config.primaryKeys显式声明逻辑主键并配合exactly_once = true,但前提是该列在源数据中确实保持唯一(详见 MySQL CDC 文档:读取没有主键的表)。

五、进阶:从"能跑"到"跑得稳"的推荐配置

5.1 低流量场景:配置 Debezium 心跳

对于低流量表,binlog 只有在发生行变更时才会推进,导致下游 checkpoint 记录不到新鲜偏移。可通过debezium块配置心跳,让 binlog 位置持续滚动(心跳表需提前在 MySQL 创建):

source { MySQL-CDC { username = "st_user_source" password = "mysqlpw" table-names = ["inventory.orders"] url = "jdbc:mysql://mysql:3306/inventory" debezium { heartbeat.interval.ms = 100 heartbeat.action.query = "INSERT INTO inventory.heartbeat (ts) VALUES (NOW())" } } }

5.2 精确一次与 Schema 演进

  • 需要端到端精确一次时,Doris 侧开启sink.enable-2pc = true,并保证sink.label-prefix全局唯一;但注意 2PC 下定时刷新与sink.buffer-size不生效,提交节奏由 checkpoint 决定。
  • 需要把 MySQL 的列变更(add/drop/rename/modify column)同步到 Doris 时,在 source 开启schema-changes.enabled = true,Doris 侧使用schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST"。仓库的 mysqlcdc_to_doris_with_schema_change.conf 即演示了schema-changes.enabled = true与 Doris sink 的组合,可作参考。

5.3 数据类型映射速查

同步前建议核对类型映射,避免隐式转换导致精度损失。MySQL 侧的映射要点(详见 MySQL CDC 数据类型映射):BIGINT UNSIGNED → DECIMAL(20,0)DECIMAL(p,s) → DECIMAL(p,s)CHAR/VARCHAR/TEXT/ENUM/JSON → STRINGDATETIME/TIMESTAMP → TIMESTAMP(s)BINARY/BLOB → BYTESTINYINT(1) → BOOLEAN(受int_type_narrowing控制)。Doris 侧的映射要点(详见 Doris 数据类型映射):BOOLEAN/TINYINT/SMALLINT/INT/BIGINT/LARGEINT对应 SeaTunnel 的整数族,DECIMAL对应DECIMAL/DOUBLE/FLOATDATE/DATETIME对应DATE/TIMESTAMPARRAY/MAP分别对应ARRAY/MAP,而HLL/BITMAP/QUANTILE_STATE/STRUCT尚不支持。


六、相关文档导航

  • MySQL CDC Source 连接器文档:完整的参数表、权限与 binlog 配置、Debezium 透传、Schema 演进示例。
  • Doris Sink 连接器文档:完整的 Sink 选项表、2PC 与删除传播、自动建表模板、307 Redirect 排查。
  • 跑第一个任务:先用 FakeSource 确认本地链路正常。
  • SeaTunnel 引擎快速开始:了解 Zeta 引擎的完整部署与运行方式。
  • 部署与连接器插件安装:发行包下载、插件安装、驱动放置。
  • 配置目录:仓库默认的插件清单(含connector-cdc-mysqlconnector-doris),可作为插件名称的权威对照。

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询