ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

FlinkCDC实战:MySQL数据实时同步到Elasticsearch的完整方案与踩坑记录

FlinkCDC实战:MySQL数据实时同步到Elasticsearch的完整方案与踩坑记录 说起来有点惭愧我第一次接到“把 MySQL 数据同步到 ES”的需求时第一反应是写一个 crontab 定时脚本每天凌晨把订单表全量导一次。这个方案上线两周就被业务方打脸了用户在 App 里改了昵称搜索结果还是旧值。那个瞬间我才想明白所谓同步不是“定期搬数据”而是“持续捕获数据库变化尽快反映到 ES”。后来我在 Canal、自写 binlog 消费和 FlinkCDC 之间反复比较最终把 MySQL 数据实时同步至 ES 的链路稳定在了 FlinkCDC 上。这篇就把整个实战过程写下来包括前期选型、环境配置、SQL 和 DataStream 两种落地方式、性能调优以及上线后踩过的那些坑希望能帮正准备做类似同步的人少走弯路。1. 为什么选 FlinkCDC先把同步方案的账算清楚1.1 同步一个表背后的真实需求表面上把 MySQL 同步到 ES 只是“数据拷贝”但业务真正需要的基本上是四件事新数据要及时出现在 ES 里老数据更新后 ES 也要跟着变MySQL 里删除的数据 ES 里不能残留同步链路不能因为任务重启就把数据弄丢或者重复导入。把这四条摆出来“定时脚本”基本被否了它只能处理第一二条删除事件难以识别延迟至少在分钟级数据量大了以后全量抽取成本还会爆炸。那 Canal 呢Canal 做增量确实成熟但它只解决“抓到 binlog 并送到 MQ”这一段后续还要自己写消费程序、自己维护 offset、自己处理重启恢复。如果团队里没有专门的平台支撑为了同步一个订单表去搭一套 MQ 加消费端运维成本偏重。FlinkCDC 的定位刚好卡在中间。它把 Debezium 的 binlog 采集能力和 Flink 的流式计算能力合在了一起既能用 SQL 直接声明一张“跟随 MySQL 变化的表”又天然支持 checkpoint、状态恢复、多表路由、数据转换。你可以把它理解成一条“会记账的流水线”binlog 是流水Flink 是会计checkpoint 是账本ES 是报表系统。1.2 不同方案的核心差异对照根据自己的场景选型直接看这张表就够了方案延迟删除事件重启恢复开发量典型适用场景定时全量脚本分钟级以上难处理无需恢复低数据量小、实时性无要求Canal MQ 自研消费秒级支持需自研中高团队已有 MQ 生态FlinkCDC 同步到 ES秒级支持checkpoint 自动恢复中低追求实时性和低运维成本选 FlinkCDC 还有一个被低估的好处全量快照和增量 binlog 之间的衔接是框架自动处理的。我们先是用scan.startup.mode initial启动任务它会先把当前表做一次全量快照快照期间产生的 binlog 也不会丢最后无缝切到增量模式。这个切换如果自己用脚本和 Canal 来写要处理“快照开始到结束这段时间的变更补偿”很容易出漏数据。2. 正式开工前MySQL 和 ES 的底子必须打好2.1 MySQL binlog 参数一个都不许错FlinkCDC 的本质是伪装成 MySQL 的一个从节点去拉 binlog所以 MySQL 端必须先把 binlog 的格式、保留时间和账号权限准备好。我这边用的 MySQL 8.0配置如下[mysqld] server-id 10324 log-bin mysql-bin binlog_format ROW binlog_row_image FULL expire_logs_days 7 max_binlog_size 256Mbinlog_format ROW是硬性要求。只有行级日志才记录每行修改前后的完整值STATEMENT 格式只记录 SQL 语句FlinkCDC 拿不到可靠的 before/after 数据。MySQL 8.0 默认就是 ROW但 5.7 不一定建议启动任务前先确认mysql -h127.0.0.1 -ucdc_user -p -e SHOW VARIABLES LIKE binlog_format;binlog_row_image FULL也很重要它决定 binlog 里是否记录整行数据。如果被改成 MINIMAL更新事件可能只记录变更列很多业务场景下会导致 ES 文档信息缺失。server-id这一项特别容易被忽略。每个连接到同一个 MySQL 实例的 CDC 任务都必须使用不同的 server-id因为 MySQL 的 binlog dump 协议会校验这个 ID。之前我见过多个任务不配置 server-id默认值撞车导致其中一个任务被 MySQL 强制断开ES 里静默缺数据。后面有专门的排查案例。expire_logs_days也别设太短。任务一旦暂停或重启需要从之前的 binlog offset 继续消费如果 binlog 被提前清理任务会因为找不到坐标直接失败。建议至少保留 7 天。2.2 给 CDC 账号授权缺一个权限都会出问题同步账号不建议直接使用 root单独建一个最小权限账号即可CREATE USER cdc_user% IDENTIFIED BY Cdc2024; GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO cdc_user%; FLUSH PRIVILEGES;这里需要解释一下为什么有RELOAD和SHOW DATABASES。RELOAD用于执行FLUSH TABLES WITH READ LOCK这是全量快照阶段保证一致性所必需的SHOW DATABASES用于发现库表。如果没有这些权限任务很可能在快照阶段报权限错误或者启动后一直卡在初始化。注意如果公司 DBA 不允许给RELOAD也可以把连接器配置里的snapshot.locking.mode调成none但这会降低快照强一致性业务能接受“快照瞬间数据接近一致但不定点”的场景才建议这么做。2.3 ES 侧先建好索引模板别让 ES 猜字段类型ES 和 MySQL 完全不一样。MySQL 写入前有强约束ES 默认是“先写入再推断”等索引里字段类型已经生成你再想改就麻烦了。所以正式同步前我建议先建一个索引模板把字段类型固下来。PUT _index_template/t_order_template { index_patterns: [t_order_index*], template: { settings: { number_of_shards: 3, number_of_replicas: 1 }, mappings: { properties: { id: { type: long }, orderNo: { type: keyword }, userId: { type: long }, amount: { type: double }, status: { type: integer }, createTime: { type: date, format: yyyy-MM-dd HH:mm:ss }, updateTime: { type: date, format: yyyy-MM-dd HH:mm:ss } } } } }订单金额之所以用 double是因为在这个同步场景里 ES 主要负责检索和展示不做精确计算精确金额仍然以 MySQL 为准。如果对金额精度要求极端严格需要用 scaled_float 或者 string 存分后文踩坑部分会再提。3. 最少代码跑通Flink SQL 直接同步单表到 ES3.1 依赖与 SQL 客户端准备我用的是 Flink 1.17.2 加 Flink CDC 2.4.0 的组合。如果用 Flink SQL 客户端直接把对应连接器 jar 放到 Flink 的 lib 目录下然后在 SQL 里声明 MySQL CDC source 和 ES sink 即可。# 示例 jar 名称具体版本以 Maven Central 搜到的为准 flink-sql-connector-mysql-cdc-2.4.0.jar flink-sql-connector-elasticsearch7-1.17.2.jarMaven 坐标的选择有一个常见误区老教程里的com.ververica:flink-connector-mysql-cdc是 2.x 早期阶段的写法后来统一到org.apache.flink下。写新任务之前去 Maven Central 搜对应版本千万别死记硬背坐标。3.2 建表语句source 表就是“会变的 MySQL 表”假设 MySQL 有张订单表t_orderCREATE TABLE t_order ( id BIGINT AUTO_INCREMENT PRIMARY KEY, order_no VARCHAR(32) NOT NULL, user_id BIGINT NOT NULL, amount DECIMAL(10, 2) NOT NULL, status TINYINT NOT NULL DEFAULT 0, create_time DATETIME NOT NULL, update_time DATETIME NOT NULL ) ENGINEInnoDB;在 Flink SQL 里这样建 CDC source 表CREATE TABLE order_cdc ( id BIGINT PRIMARY KEY, order_no STRING, user_id BIGINT, amount DECIMAL(10, 2), status INT, create_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname 192.168.1.10, port 3306, username cdc_user, password Cdc2024, database-name shop, table-name t_order, scan.startup.mode initial, server-time-zone Asia/Shanghai );再建目标索引的 sink 表CREATE TABLE order_es ( id BIGINT PRIMARY KEY, order_no STRING, user_id BIGINT, amount DECIMAL(10, 2), status INT, create_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector elasticsearch-7, hosts http://192.168.1.20:9200, index t_order_index, sink.bulk-flush.max-actions 1000, sink.bulk-flush.interval 1000, format json );最后一条 SQL 提交任务INSERT INTO order_es SELECT id, order_no, user_id, amount, status, create_time FROM order_cdc;scan.startup.mode initial的含义是任务启动后先扫描 MySQL 当前全量数据写入 ES再接着消费 binlog 增量。如果你只需要从当前时刻开始同步新数据可以改成latest-offset但初始化空索引、需要补历史数据时一定要用 initial。3.3 同步结果验证插入、更新、删除都要测任务跑起来以后先做三个最基础的验证。第一MySQL 里插入一条订单INSERT INTO t_order (order_no, user_id, amount, status, create_time, update_time) VALUES (SO20240613001, 88, 199.00, 0, NOW(), NOW());等一两秒到 ES 里查curl -s http://192.168.1.20:9200/t_order_index/_search?qorderNo:SO20240613001第二更新这条订单的status确认 ES 文档同步变成新值。第三删除这条记录确认 ES 里对应文档也被删除。这个环节建议做成自动化脚本以后每次改造同步链路跑一遍就知道有没有回归问题。4. 单表搞不定的场景用 DataStream 做多表路由和定制清洗4.1 为什么还需要 Java APIFlink SQL 适合单表、字段一一对应的同步。业务一旦出现这几种情况SQL 就不够灵活了多张 MySQL 表要进同一个 ES 索引比如订单宽表需要按来源表拼不同字段某个字段在 ES 里要特殊处理比如把状态码映射成可读文案表名或库名不同但内容要一起检索需要根据source.table动态路由到不同索引。DataStream API 就是用来处理这些“必须写代码”的场景的。它的开发量确实比 SQL 大一点但换来的是完全可控的事件处理逻辑。4.2 单任务消费多表按来源动态写索引Maven 依赖大致如下版本按你的 Flink 去匹配dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-mysql-cdc/artifactId version2.4.0/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-elasticsearch7/artifactId version1.17.2/version /dependency下面是一个简化但能跑通的核心骨架public class MysqlMultiTableToEs { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(5000); env.getCheckpointConfig().setCheckpointStorage(file:///data/flink-cdc-checkpoint); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(2000); MySqlSourceString source MySqlSource.Stringbuilder() .hostname(192.168.1.10) .port(3306) .databaseList(shop) .tableList(shop.t_order, shop.t_user) .username(cdc_user) .password(Cdc2024) .serverTimeZone(Asia/Shanghai) .deserializer(new JsonDebeziumDeserializationSchema()) .startupOptions(StartupOptions.initial()) .build(); DataStreamSourceString stream env.fromSource( source, WatermarkStrategy.noWatermarks(), mysql-cdc-multi-table); stream.map(new RouteAndFormatFunction()) .filter(Objects::nonNull) .addSink(createEsSink()); env.execute(mysql-cdc-sync-to-es); } }RouteAndFormatFunction 里做的就是把 Debezium 的原始 JSON 转成 ES 的 IndexRequest 或 DeleteRequestpublic static class RouteAndFormatFunction implements MapFunctionString, Object { Override public Object map(String raw) throws Exception { JSONObject record JSON.parseObject(raw); String op record.getString(op); JSONObject source record.getJSONObject(source); String table source null ? : source.getString(table); if (d.equals(op)) { JSONObject before record.getJSONObject(before); StringBuilder docId new StringBuilder(); docId.append(before.getString(id)); return new DeleteRequest(indexFor(table), docId.toString()); } JSONObject after record.getJSONObject(after); if (after null) { return null; } IndexRequest request new IndexRequest(indexFor(table)) .id(after.getString(id)) .source(buildEsDoc(table, after).toJSONString(), XContentType.JSON); return request; } private String indexFor(String table) { if (t_order.equals(table)) return t_order_index; if (t_user.equals(table)) return t_user_index; return null; } }这里的核心思想和 SQL 接触的一样MySQL 主键id直接映射成 ES 文档_id。这样 update 事件到达时ES 会直接覆盖同 ID 文档天然具备 upsert 语义delete 事件也能精确删除对应文档。4.3 同一个 MySQL 实例多个库、多个表怎么处理如果只需要一个 MySQL 实例里的多张表一个MySqlSource就够tableList里用逗号分隔多个表名。不同业务表写不同索引时在代码里根据source.table路由即可。如果涉及多个 MySQL 实例比如主库和从库各有一份数据可以创建多个MySqlSource再通过connect或union合并数据流。这样的好处是每个源可以独立设置 server-id、账号和表清单但要注意每个源都必须有唯一的 server-id后面踩坑部分我会展开。5. 性能与可靠性把同步任务当成正式生产任务对待5.1 Checkpoint 是 CDC 任务的保命符FlinkCDC 之所以比 Canal 加自研消费更可靠很大的原因是 checkpoint 自动保存了消费位点。Flink 的 state 里存着 Debezium 当前的 binlog offset任务重启后可以从 offset 恢复保证“不丢数据、不重复消费”的精确一次语义。如果不开 checkpoint任务一重启可能从latest-offset接着消费也可能重新走一遍全量初始化结果是增量数据丢一段或者 ES 被重复写入大量数据。生产环境我习惯这样配env.enableCheckpointing(5000); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(2000); env.getCheckpointConfig().setCheckpointTimeout(60000); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3, Time.seconds(10)));setMinPauseBetweenCheckpoints很重要它防止两个 checkpoint 之间挤得太紧影响正常数据处理。checkpoint 存储尽量不要用本地磁盘尤其是集群部署时建议指向 HDFS 或对象存储否则 TaskManager 重启后状态就找不到了。5.2 写入 ES 的批量参数不要用默认值硬扛ES sink 的批量参数直接影响吞吐。默认参数不是按“同步大表”的场景调的我常用的配置如下sink.bulk-flush.max-actions 1000, sink.bulk-flush.max-size 5mb, sink.bulk-flush.interval 2000这几个参数合起来的含义是攒够 1000 条、或者 5MB、或者 2 秒就批量刷一次。数值太小请求频繁ES 压力大数值太大队列积压延迟变高。同步场景的体量下上面这组值是一个比较稳的起点再根据实际吞吐微调。FlinkCDC 的 source 端也不是无脑加并行度就能提速的。MySQL CDC 的 binlog 读取受 Debezium 设计限制单表事件往往是单并行度处理真正的吞吐瓶颈多半在 ES 写入端。所以优化方向应该是提高 sink 并行度而不是盲目给 source 加并行。5.3 全量初始化时ES 侧要做针对性调整scan.startup.mode initial启动后会先做全量快照。此时 ES 如果保持生产状态的number_of_replicas 1、refresh_interval 1s每个分片都要复制副本、每秒刷新写放大很严重。我一般在全量阶段临时调低这些参数PUT t_order_index/_settings { index: { number_of_replicas: 0, refresh_interval: 30s } }等到 binlog 增量追平、文档数和 MySQL 基本一致后再改回来PUT t_order_index/_settings { index: { number_of_replicas: 1, refresh_interval: 1s } }这个操作看起来不起眼但在同步几百 GB 级别的订单表时能把初始化时间缩短一倍以上。5.4 监控指标lag 比文档数量更值得关注任务上线后判断“是否同步成功”不能只盯着 ES 文档数。文档数只是全量结果增量链路是否健康要看两个指标消费延迟和当前心跳事件时间。在 Flink Web UI 的指标面板里可以找到 Debezium 暴露的currentEmitEventTimeLag和millisBehindSource。millisBehindSource表示当前消费到的 binlog 时间与 MySQL 最新 binlog 时间的差。这个值如果长期高企说明 source 或 sink 有瓶颈如果它稳步回落说明任务正在追进度。每次上线新同步任务我的第一步都是看这个指标能不能稳定到几秒以内而不是只查 ES 有没有数据。6. 踩坑实录五个真实问题的排查链路和最终修复6.1 server-id 冲突为什么一个任务突然断流ES 还静默缺数据现象是某天线上同时跑两个 MySQL 同步任务其中一个索引的文档数在半小时内不涨反降。查 Flink 日志看到了类似 “A slave with the same server_uuid/server_id as this slave has connected to the master” 的报错。排查链路是这样的先确认两个任务连接的是同一个 MySQL 实例再看它们的 source 配置发现都没有显式配置server-id。Debezium 默认生成的 server-id 可能相同两个 CDC 任务同时伪装成同一个从库MySQL 就会把后连接的连接踢掉导致其中一个任务持续“拉不到新事件”。修复方案很简单每个 MySqlSource 都配置不同的 server-id 或一个范围。MySqlSource.Stringbuilder() .serverId(5400-5500) ...这里给的是一个范围不是单个数字因为 CDC 在并发读取多个表时可能需要多个通道 ID。SQL 方式同理在 WITH 参数里加server-id 5400-5500。注意不同任务之间 server-id 的取值区间不能重叠否则问题还会复现。6.2 MySQL 5.7 默认 binlog_format 不是 ROW任务启动成功但数据不对另一台 MySQL 5.7 测试环境FlinkCDC 任务启动时没有任何报错插入数据也能同步到 ES但某次更新用户昵称后ES 里的昵称一直不变。排查过程是先看 source 表格式发现binlog_format是 STATEMENT。为什么 STATEMENT 下部分更新也能同步因为 Debezium 在某些情况下仍然可以从语句结果推导部分变更但涉及多行、复杂条件更新时它能拿到的 before/after 就是不完整的。这个坑最恶心的地方在于“时好时坏”容易让人误以为是程序逻辑问题。修复方法是修改 MySQL 配置并重启确认binlog_formatROW、binlog_row_imageFULL。改完后再做一遍全量 update 验证。6.3 全量快照阶段 TaskManager OOM第一次同步刚上线全量阶段跑了没几分钟TaskManager 就内存溢出作业直接失败。日志里是大片 GC 和java.lang.OutOfMemoryError: Java heap space。原因是全量快照时Debezium 会分批读取大表数据到内存。默认的snapshot.fetch.size对超大表来说太大一次拉取的数据量加上快照缓存把 TaskManager 堆内存顶爆了。修复分两步。第一步调小单批拉取行数Properties debeziumProps new Properties(); debeziumProps.setProperty(snapshot.fetch.size, 100); MySqlSource.Stringbuilder() .debeziumProperties(debeziumProps) ...第二步给 TaskManager 留足内存避免和其他流任务混跑。别指望调小 fetch size 就万事大吉全量扫描是一连串查询内存和耗时需要一起权衡。6.4 时区问题ES 里时间总是差八小时同步订单表后对比 ES 文档的createTime发现比 MySQL 里的原始时间整整晚了 8 小时。这个问题的根因是 JDBC 连接 MySQL 时驱动会把 DATETIME 字段按照一个默认时区解析而 MySQL 服务端是 UTC本地是东八区最终在同步链路里被多做了一次时区转换。修复方式是统一时区最简单的是在 source 表 WITH 参数里显式指定server-time-zone Asia/Shanghai同时在 JDBC URL 层面也不要再画蛇添足加useTimezone之类的参数。时区问题排查起来很费劲因为不是所有字段都会错只有日期时间字段会被“悄悄”改掉建议任何同步任务开局就把时区写死。6.5 ES mapping 被动态映射带偏搜索语义不对还有一次是金额字段在 ES 里被映射成了 string。原因是 MySQL 里DECIMAL经过 Flink 类型的默认转换在写入 JSON 时变成了带引号的字符串ES 动态映射就把字段判断成 text。这种问题在写入量小的时候很难发现等数据量大了再改 mapping 就要 reindex。正确做法是同步前先把索引模板建好金额字段明确写成double或scaled_float。真的要改 mapping 时应该通过创建一个新索引再用 reindex 迁移数据不要在旧索引上做永久性修改。最后说一句实在话这套链路跑通之后我养成了一个习惯每次改动同步逻辑都会在 MySQL 里故意做一次整表更新比如全表把status字段统一刷新一遍然后用 ES 的查询统计文档数和更新时间分布确认没有漏事件也没多事件。这个动作看似笨但已经帮我抓出过至少三次 server-id 冲突和一次权限不足导致的慢消费。如果你也准备把 FlinkCDC 同步 MySQL 到 ES 这条链路用在生产环境建议把这一条写进发布清单它会救你于半夜的告警电话。
返回列表