ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Apache SeaTunnel 实战:MySQL CDC 实时同步到 Doris 的完整配置与验证指南

2026/9/17 19:26:34 拓冰建站 浏览量
Apache SeaTunnel 实战:MySQL CDC 实时同步到 Doris 的完整配置与验证指南 Apache SeaTunnel 实战MySQL CDC 实时同步到 Doris 的完整配置与验证指南【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本指南围绕 SeaTunnel 官方 Recipedocs/en/getting-started/recipes/mysql-cdc-to-doris.md展开完整讲解如何用 SeaTunnel 从 MySQL 捕获行级变更CDC并持续写入 Doris覆盖插件安装、binlog 与权限准备、HOCON 配置编写、任务启动与结果验证并补充仓库源码与端到端测试证据。读完本文你将能够独立搭建一条「MySQL 全量 增量 → Doris」的实时同步管道并掌握server-id、sink.label-prefix、删除传播等关键参数的正确用法。场景概述为什么要用 MySQL CDC 同步到 Doris当业务需要把 MySQL 中的订单等核心表实时同步到 Doris分析型数仓时最常见的诉求是启动时先做一次全量快照之后持续接收 MySQL binlog 中的增量变更INSERT / UPDATE / DELETE 三种变更都要原样传导到 DorisDoris 侧无需人工预先建表表结构可以由 MySQL 主键元数据自动推导生成。SeaTunnel 的connector-cdc-mysql与connector-doris正好覆盖这条链路前者基于快照 binlog 增量实现变更捕获后者通过 Doris Stream Load 批量导入并支持唯一键模型的删除传播。本文以官方 Recipe 为骨架给出可复制的完整方案。前置条件环境、插件与 JDBC 驱动1. 本地可运行的 SeaTunnel 环境先确认你已经完成首次任务并能在本地跑通参考 run-your-first-job。本地环境要求安装 Java8 或 11 及以上并配置好JAVA_HOME详见 部署文档。2. 只安装本 Recipe 需要的两个插件从 2.2.0-beta 起SeaTunnel 二进制包默认不再携带所有 connector需要先用install-plugin.sh按需安装。编辑${SEATUNNEL_HOME}/config/plugin_config只保留如下内容--seatunnel-connectors-- connector-cdc-mysql connector-doris --end--然后执行安装脚本并确认两个 connector 的 jar 已就位cd ${SEATUNNEL_HOME} sh bin/install-plugin.sh ls connectors | rg connector-(cdc-mysql|doris)关于plugin_config的完整语义可参考仓库中的 config/plugin_config该文件列出了所有可选 connector--connectors-v2--段是连接器插件的选择清单install-plugin.sh会依据它下载对应 jar。仓库根目录的 plugin-mapping.properties 则维护了插件名与 Maven 制品artifactId的映射关系是插件解析与加载的依据。3. 准备 MySQL JDBC 驱动如果使用 SeaTunnel Zeta默认引擎将mysql-connector-java如 8.0.28放入${SEATUNNEL_HOME}/lib并确认可见ls ${SEATUNNEL_HOME}/lib | rg mysql-connector如果改用 Flink 或 Spark 引擎则需要把同一驱动 jar 放到对应引擎的插件目录${SEATUNNEL_HOME}/plugins/详见 MySQL CDC 源文档 的 Using Dependency 一节——这正是 CDC 读取 MySQL 元数据、执行快照查询的 JDBC 依赖。4. 准备 MySQL 源表必须有稳定主键CDC 的 UPDATE / DELETE 事件需要依赖行标识primary key才能在下游正确回放因此源表必须声明主键。Recipe 中给出的建表语句如下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());从源码角度看主键的作用不仅体现在下游回放在connector-cdc-mysql的源码中快照阶段会把大表按主键拆分成多个 split 并行读取对应MySqlChunkSplitter位于 MySqlChunkSplitter.java无主键表则无法可靠拆分也不具备确定性的行身份。5. 创建 CDC 专用账号并授权MySQL CDC 基于 Debezium 嵌入式引擎实现需要与 MySQL CDC 源 相同的权限SELECT读取表数据与元数据、RELOAD用于快照加锁的一致性读取、SHOW DATABASES、以及REPLICATION SLAVE/REPLICATION CLIENT订阅 binlogCREATE 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;6. 确认 binlog 已开启且为 ROW 格式MySQL CDC 依赖 binlog 记录行级变更必须先确认相关变量SHOW VARIABLES WHERE variable_name IN (log_bin, binlog_format, binlog_row_image);期望值log_bin ON、binlog_format ROW、binlog_row_image FULL。若未开启在my.cnf中追加如下配置并重启 MySQL[mysqld] server-id 223344 log_bin mysql-bin binlog_format ROW binlog_row_image FULL补充说明binlog_format ROW保证 binlog 记录的是行级前后镜像而非 SQL 语句是 CDC 捕获 UPDATE 前后值与 DELETE 行的前提binlog_row_image FULLMySQL 5.6 要求确保每一行都记录完整列镜像。如果需要更精确的断点续传还可以按 MySQL CDC 源文档 的建议开启 GTIDgtid_mode on、enforce_gtid_consistency on。7. 准备 Doris 目标库Recipe 中schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST因此首次启动时 SeaTunnel 会根据 MySQL 主键元数据在 Doris 自动创建sync_demo.orders。这里只需预先建库CREATE DATABASE IF NOT EXISTS sync_demo;最小配置一条可运行的 CDC 管道将以下配置保存为config/mysql-cdc-to-doris.confenv { 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 , } } }env 段说明job.mode STREAMING声明为流式任务任务不会自行结束持续消费 binlogcheckpoint.interval 5000每 5 秒做一次 checkpoint。Doris sink 的提交行为与 checkpoint 强相关——当数据量未达到doris.batch.size阈值时只有 checkpoint 会触发一次 Stream Load 提交详见 Doris sink 文档 的 Tuning Guide因此该值同时决定了低流量场景下的最大可见延迟parallelism 1本示例单并行度即可。sourceMySQL-CDC段说明参数示例值说明plugin_outputorders_cdc为数据流命名sink 侧用plugin_input与之对接也可不写按默认顺序配对startup.modeinitial先做全量快照再无缝切换为 binlog 增量读取是 Recipe 的默认推荐server-id5652CDC 读取器在 MySQL 侧的复制 ID必须与集群内其他 replica / CDC 任务不冲突多并行或多表读取时可配置区间如5400-5408table-names[inventory.orders]监控的表清单格式为库名.表名urljdbc:mysql://mysql:3306/inventoryJDBC 连接串mysql为容器/主机名可按实际环境替换为 IP关于startup.mode源码 MySqlIncrementalSourceOptions.java 中定义了完整枚举initial、earliest、latest、specific、timestamp默认值为INITIALstop.mode默认NEVER。initial的工作方式为启动时对目标表做一致性快照快照完成后自动从快照起始时记录的 binlog 位置继续消费切换过程不丢事件。其他有用的启动模式包括latest跳过历史快照只从当前最新 offset 开始消费历史数据不需要时可大幅加速启动specific配合startup.specific-offset.file/startup.specific-offset.pos从指定 binlog 文件与位置开始timestamp配合startup.timestamp从指定时间戳Unix 毫秒开始。sinkDoris段说明参数示例值说明fenodesdoris-fe:8030Doris FE 的 HTTP 地址fe_ip:fe_http_portdatabase/tablesync_demo/orders目标库表多表场景可用${database_name}/${table_name}占位符继承上游库表名sink.label-prefixorders-cdcStream Load 的 label 前缀多个并发任务间必须全局唯一否则会触发 Doris Label already exists 冲突sink.enable-deletetrue开启删除传播CDC 源产生的 DELETE 会转成 Doris 的批量删除标记要求目标表为Unique Key 模型schema_save_modeCREATE_SCHEMA_WHEN_NOT_EXIST表不存在时自动建表已存在则跳过doris.configformat csv、column_separator ,Stream Load 的 data_desc 参数CSV 模式需显式指定分隔符也可以改用format jsonread_json_by_line trueschema_save_mode的完整取值见 Doris sink 文档为RECREATE_SCHEMA不存在则创建存在则删除重建、CREATE_SCHEMA_WHEN_NOT_EXIST推荐存在即跳过、ERROR_WHEN_SCHEMA_NOT_EXIST不存在则报错、IGNORE不处理。自动建表使用save_mode_create_template模板默认生成ENGINEOLAP、UNIQUE KEY (主键)、DISTRIBUTED BY HASH (主键)的建表语句占位符包括${database}、${table}、${rowtype_fields}、${rowtype_primary_key}等可自定义模板来调整分桶、副本等属性。从实现上看Doris sink 采用「缓存 批量 Stream Load」的写入路径源码中 DorisSinkOptions.java 定义了sink.buffer-size默认 256 * 1024、sink.buffer-count默认 3、sink.max-retries默认 3、sink.check-interval默认 10000等参数DorisStreamLoad.java 实现了实际的 Stream Load HTTP 请求与 2PC 提交/中止控制。若需严格 exactly-once可设置sink.enable-2pc true此时注意sink.buffer-size将不再生效提交只由 checkpoint 触发。运行任务在${SEATUNNEL_HOME}下以本地模式启动cd ${SEATUNNEL_HOME} ./bin/seatunnel.sh --config ./config/mysql-cdc-to-doris.conf -m local因为这是流式 CDC 管道任务会持续运行请在任务运行期间执行下面的验证 SQL而不是等任务退出。验证结果全量 增量变更如何传导到 Doris验证分三步覆盖了「全量快照 → 增量 INSERT → 增量 UPDATE → 增量 DELETE」四个阶段启动任务等待初始快照完成此时1001、1002两行应已进入 Doris在 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;在 Doris 侧查询最终状态SELECT COUNT(*) FROM sync_demo.orders; SELECT id, order_status, amount FROM sync_demo.orders ORDER BY id;期望结果能看到1001的状态为PAID、1003的状态为CREATED且1002已不存在。如果 INSERT / UPDATE / DELETE 三种变更都如实反映到了 Doris说明管道工作正常——sink.enable-delete true正是 DELETE 能被传导的关键开关。这一套验证流程在仓库的端到端测试中有对应实现connector-cdc-mysql-e2e模块下的 MysqlCDCIT.java 通过 Testcontainers 拉起 MySQL 8.0 容器并执行完整的快照 增量断言MysqlCDCWithBinlogDeleteIT.java 则专门验证 binlog 中 DELETE 事件的捕获与传导。这些测试即为本文验证步骤的自动化版本可作为排查问题时的对照参考。常见陷阱与排查清单官方 Recipe 列举了以下高频问题这里逐一给出成因与对策MySQL binlog 未开启或不是 ROW 格式CDC 无数据来源。用前文SHOW VARIABLES校验修改my.cnf后必须重启 MySQL。CDC 账号缺少复制权限表现为任务启动后无法订阅 binlog。确认账号具备REPLICATION SLAVE, REPLICATION CLIENT, RELOAD, SHOW DATABASES, SELECT权限。server-id与 MySQL 其他 replica 或其他 CDC 任务冲突MySQL 会强制断开重复 ID 的客户端。每个 CDC 任务使用唯一 ID 或不相交的 ID 区间例如任务 A 用5400-5600、任务 B 用5601-5800。sink.label-prefix被多个运行中任务复用Doris Stream Load 以 label 去重冲突会导致导入失败。每个任务使用独立前缀或在前缀中带上时间戳/随机串若历史 label 残留可在 Doris 执行CANCEL LOAD WHERE LABEL LIKE your-prefix%清理未提交事务。开启删除传播但 Doris 表模型不支持删除sink.enable-delete true要求目标表为Unique Key 模型且 Doris 版本支持批量删除0.15 默认开启Recipe 的自动建表模板默认生成 UNIQUE KEY 表恰好满足要求。源表没有稳定主键Doris 自动建表与下游 upsert 行为将不确定。如果源表确实无主键可参考 MySQL CDC 源文档 用table-names-config.primaryKeys显式指定唯一列作为行身份纯追加型负载无 UPDATE/DELETE则可保持exactly_once false不做主键声明。进阶扩展方向多表同步table-names支持多个表Doris sink 可用database ${database_name}_test、table ${table_name}_test做库表映射Schema 变更同步在 source 段开启schema-changes.enabled true当前支持add column、drop column、rename column、modify column四类 DDL 的传播MySQL CDC 源文档 中还有schema-changes.include/schema-changes.exclude过滤机制Doris sink 侧也提供了对应的 CDC 建表示例直接写入 BE 节点如遇 FE 307 重定向问题可配置benodes并设置direct_to_be true绕过 FE 转发路径2PC 时数据写入走 BE、提交/中止控制仍走 FE详见 Doris sink 文档 的 Redirect Behavior引擎选择connector-cdc-mysql支持 SeaTunnel Zeta 与 Flink 引擎Spark 暂不支持 CDC仓库 e2e 测试对 Spark 做了DisabledOnContainer标注Zeta 引擎的快速上手见 quick-start-seatunnel-engine。相关文档MySQL CDC 源文档权限、binlog、数据类型映射与全部 source 参数Doris sink 文档Stream Load、2PC、schema_save_mode、模板化建表与调优指南SeaTunnel Engine 快速上手部署与插件安装plugin_config与install-plugin.sh的完整用法【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考