1. 项目概述:MySQL到ES的数据同步实战
去年接手的一个电商项目让我深刻体会到实时数据同步的重要性。当时商品信息在MySQL中更新后,需要长达15分钟才能在搜索系统中生效,直接影响了促销活动的效果。为了解决这个问题,我选择了FlinkCDC作为数据同步方案,最终将延迟控制在毫秒级别。
FlinkCDC是Apache Flink社区基于变更数据捕获(CDC)技术开发的组件,能够实时捕获数据库的变更事件。相比传统的定时轮询或双写方案,它具有以下不可替代的优势:
- 低延迟:基于数据库日志解析,变更事件产生后立即处理
- 低负载:不依赖查询业务表,对源库压力极小
- 一致性:保证至少一次(at-least-once)的事件投递语义
- 全量+增量:支持历史数据初始化与实时变更同步的统一处理
典型应用场景包括:
- 搜索索引构建(如本文的MySQL→ES场景)
- 数据仓库实时ETL
- 多级缓存一致性维护
- 微服务间的数据依赖解耦
2. 环境准备与组件配置
2.1 基础环境搭建
建议使用以下版本组合以避免兼容性问题:
# 组件版本 Flink 1.15.3 Flink CDC Connectors 2.3.0 MySQL 5.7+ (需开启binlog) Elasticsearch 7.10+MySQL必须开启binlog并配置ROW模式:
# 检查当前配置 SHOW VARIABLES LIKE 'log_bin'; SHOW VARIABLES LIKE 'binlog_format'; # 修改my.cnf [mysqld] server-id = 1 log_bin = mysql-bin binlog_format = ROW binlog_row_image = FULL expire_logs_days = 7注意:生产环境建议为FlinkCDC创建专用账号并授权:
CREATE USER 'flinkcdc'@'%' IDENTIFIED BY 'SecurePwd123!'; GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'flinkcdc'@'%'; FLUSH PRIVILEGES;
2.2 Elasticsearch索引设计
以电商商品表为例,合理的ES索引映射应该考虑以下因素:
PUT /products { "settings": { "number_of_shards": 3, "number_of_replicas": 1, "refresh_interval": "30s" }, "mappings": { "properties": { "id": {"type": "keyword"}, "name": { "type": "text", "analyzer": "ik_max_word", "fields": {"raw": {"type": "keyword"}} }, "price": {"type": "scaled_float", "scaling_factor": 100}, "stock": {"type": "integer"}, "categories": {"type": "keyword"}, "attributes": { "type": "nested", "properties": { "name": {"type": "keyword"}, "value": {"type": "keyword"} } }, "create_time": { "type": "date", "format": "yyyy-MM-dd HH:mm:ss||epoch_millis" } } } }3. 核心实现解析
3.1 FlinkCDC数据捕获配置
使用Java API创建MySQL CDC源:
DebeziumSourceFunction<String> sourceFunction = MySQLSource.<String>builder() .hostname("mysql-host") .port(3306) .username("flinkcdc") .password("SecurePwd123!") .databaseList("ecommerce") .tableList("ecommerce.products") .deserializer(new JsonDebeziumDeserializationSchema()) .startupOptions(StartupOptions.initial()) .serverTimeZone("Asia/Shanghai") .build();关键参数说明:
startupOptions:支持多种初始化模式initial:先做全量快照,然后接增量latest:仅从当前开始消费增量timestamp:从指定时间点开始
serverTimeZone:必须与MySQL服务器时区一致deserializer:控制事件解析格式
3.2 数据转换与写入ES
构建Elasticsearch Sink:
List<HttpHost> httpHosts = Arrays.asList( new HttpHost("es-node1", 9200, "http"), new HttpHost("es-node2", 9200, "http") ); ElasticsearchSink.Builder<String> esSinkBuilder = new ElasticsearchSink.Builder<>( httpHosts, (element, ctx, indexer) -> { // 解析CDC事件 JsonNode jsonNode = JsonUtils.parse(element); String op = jsonNode.get("op").asText(); // 只处理插入/更新事件 if ("c".equals(op) || "u".equals(op)) { IndexRequest request = Requests.indexRequest() .index("products") .id(jsonNode.get("after").get("id").asText()) .source(element); indexer.add(request); } else if ("d".equals(op)) { DeleteRequest request = Requests.deleteRequest("products") .id(jsonNode.get("before").get("id").asText()); indexer.add(request); } } ); // 批量写入配置 esSinkBuilder.setBulkFlushMaxActions(1000); esSinkBuilder.setBulkFlushInterval(5000); esSinkBuilder.setBulkFlushBackoff(true); esSinkBuilder.setBulkFlushBackoffType(BackoffType.EXPONENTIAL); esSinkBuilder.setBulkFlushBackoffDelay(3000); esSinkBuilder.setBulkFlushBackoffRetries(3);3.3 完整作业组装
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 启用检查点(保证Exactly-Once语义) env.enableCheckpointing(30000); env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(500); env.getCheckpointConfig().setCheckpointTimeout(60000); // 构建Pipeline DataStreamSource<String> source = env.addSource(sourceFunction); source.addSink(esSinkBuilder.build()); // 执行作业 env.execute("MySQL-to-ES-Sync");4. 生产环境优化实践
4.1 性能调优参数
// 并行度设置(根据CPU核心数调整) env.setParallelism(4); // 网络缓冲区优化 env.setBufferTimeout(100); env.getConfig().setNettyShuffleMode(NettyShuffleMode.ALL_EDGES_BLOCKING); // 状态后端配置 env.setStateBackend(new RocksDBStateBackend("hdfs://namenode:8020/flink/checkpoints", true));4.2 容错与监控
- 断点续传:通过保存的checkpoint恢复作业
# 从检查点重启 bin/flink run -s hdfs://namenode:8020/flink/checkpoints/savepoint-123456 \ -c com.etl.Main your-job.jar- 监控指标:通过Prometheus采集关键指标
# flink-conf.yaml配置 metrics.reporter.prom.class: org.apache.flink.metrics.prometheus.PrometheusReporter metrics.reporter.prom.port: 9250-9260- 告警规则(示例):
- 消费延迟 > 10s
- 检查点失败率 > 20%
- ES批量写入错误率 > 5%
4.3 常见问题排查
问题1:CDC连接中断
- 现象:作业持续重启,日志显示"Connection reset by peer"
- 解决方案:
- 增加心跳间隔:
heartbeat.interval.ms=60000 - 配置连接池:
connection.pool.size=5
- 增加心跳间隔:
问题2:ES写入瓶颈
- 现象:sink延迟持续增长
- 优化方向:
- 增加ES索引刷新间隔:
"refresh_interval": "60s" - 调整批量参数:
setBulkFlushMaxActions(5000) - 升级ES集群硬件或增加节点
- 增加ES索引刷新间隔:
问题3:数据不一致
- 排查步骤:
- 检查binlog位置是否正常推进
- 验证Flink检查点是否完整
- 对比MySQL与ES的count记录数差异
5. 进阶应用场景
5.1 多表关联同步
通过Lookup Join实现维度表关联:
// 主表CDC源 DataStream<JsonNode> orders = env.addSource(orderSource); // 维度表CDC源 DataStream<JsonNode> users = env.addSource(userSource); // 构建时态表 TemporalTableFunction userTable = users .keyBy(node -> node.get("id").asText()) .createTemporalTableFunction( node -> Instant.parse(node.get("update_time").asText()), "user_info"); // 注册函数 env.registerFunction("userInfo", userTable); // 执行关联查询 DataStream<EnrichedOrder> result = orders .keyBy(node -> node.get("user_id").asText()) .process(new TemporalJoinProcessFunction());5.2 数据结构变更处理
应对MySQL DDL变更的策略:
- Schema Evolution:使用Avro格式存储schema历史版本
- 双跑过渡:新旧schema程序并行运行直至数据迁移完成
- 离线补偿:通过全量扫描修复不一致数据
5.3 数据清洗与转换
典型ETL处理链示例:
source // 过滤无效数据 .filter(node -> !node.get("id").isNull()) // 字段脱敏 .map(node -> { JsonNode cloned = node.deepCopy(); ((ObjectNode)cloned).put("phone", maskPhone(node.get("phone").asText())); return cloned; }) // 扁平化嵌套结构 .flatMap(new NestedFieldFlattener()) // 写入ES .addSink(esSink);在实际项目中,这套方案将百万级商品数据的同步延迟从原来的15分钟降低到500毫秒以内。值得注意的是,ES的写入性能与索引设计密切相关,建议在正式上线前进行充分的压力测试。我曾遇到过一个案例:由于未合理设置分片数,写入吞吐量始终上不去,后来通过重建索引将性能提升了3倍。