ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

FlinkCDC同步性能卡死?读写解耦+Kafka并行度优化实战

FlinkCDC同步性能卡死?读写解耦+Kafka并行度优化实战 先说结论如果你的 FlinkCDC 数据同步任务遇到同步性能无法提升、怎么调 Sink 并行度吞吐都纹丝不动的情况大概率问题不在 Sink 端而是整条链路的写入并行度被上游 Source 的单通道给锁死了。这个坑我踩了一整天才彻底定位当时从发现写入 Doris 的吞吐上不去到最后通过读写解耦方案把性能提升了接近 8 倍中间那一整条排查链路我认为很值得单独写一篇复盘。不管你现在用的是 FlinkCDC 2.x 还是 3.x只要架构还是“MySQL CDC 直连目标库”这篇文章的内容就都适用其中包含可照抄的参数配置、反压判定技巧还有几个特别容易误判的细节。1. 现象还原Sink 并行度调到 8吞吐反而没变化先交代一下任务背景。我当时要同步的是线上一个订单库MySQL 单实例存量订单表接近 800 万行日均新增 50 万左右。目标库是 Doris用来做实时 OLAP 报表。Flink 版本 1.16.2FlinkCDC 用的 2.3.0代码基于 DataStream API没有走 SQL Client这样方便后面控制并行度和分区逻辑。任务的初始拓扑非常简单MySQL binlog - FlinkCDC Source - Doris SinkSource 端按照官方推荐的方式配置用MySqlSource的 Builder 构建启动模式用的initial()因为第一次要同时处理全量存量数据和增量日志。Sink 端用的是 Doris Connector为了拉高写入并行度我显式指定了.setParallelism(8)。DataStreamString source env.addSource( MySqlSource.Stringbuilder() .hostname(mysql-primary) .port(3306) .databaseList(shop) .tableList(shop.t_order) .username(cdc_user) .password(******) .deserializer(new JsonDebeziumDeserializationSchema()) .startupOptions(StartupOptions.initial()) .serverTimeZone(Asia/Shanghai) .build() ); Properties streamLoadProps new Properties(); streamLoadProps.setProperty(format, json); streamLoadProps.setProperty(read_json_by_line, true); source.addSink( DorisSink.builder() .setDorisOptions(DorisOptions.builder() .setFenodes(doris-fe:8030) .setTableIdentifier(olap.t_order) .setUsername(doris_user) .setPassword(******) .build()) .setDorisExecutionOptions(DorisExecutionOptions.builder() .setLabelPrefix(doris-cdc) .setStreamLoadProp(streamLoadProps) .build()) .build() ).setParallelism(8);因为开会讨论时大家预期的是“8 个并行度写入怎么也能跑到几万行每秒”所以任务上线后我特意做了一轮压测。结果很尴尬。Sink 并行度批刷新 rows实测吞吐现象1默认 100约 8000 行/s反压 High4默认 100约 7600 行/s反压 High8默认 100约 7900 行/s反压 High85000约 8100 行/s无实质改善82000约 8000 行/s无实质改善把 Sink 并行度从 1 调到 8吞吐不仅没涨反而还出现轻微下降。这里有一个反直觉的细节并行度调高之后Doris 端需要处理更多并发的 stream load 请求如果连接池和标签管理跟不上反而会因为频繁的 label 冲突和连接争抢拖慢单次写入。所以最开始我把问题定位在 Doris Sink 参数上调了 batch 大小、刷新间隔、连接池但结果基本没有变化。真正让我警觉的是一个数据不管怎么调Source 端每秒产出的记录数始终稳定在 8000 行左右像是有人在源头把水龙头拧死了一样。当时我判断问题大概率不在 Sink而在更上游的地方。2. 排查链路从数据倾斜一路追溯到 Source 单通道定位这类同步性能问题最忌讳一上来就盲改参数。我按下面的顺序一步步收窄范围最后抓到了根因。2.1 第一步确认 Sink 并行度真的生效了吗先看 Flink Web UI 的算子状态。Sink 算子的并行度确实显示为 8/8说明.setParallelism(8)是生效的。但接着看 TaskManager 的 CPU 和线程活动情况时发现只有 1 个 subtask 在持续干活其余 7 个 subtask 基本处于空闲等待状态。再点开 Sink 算子每个 subtask 的接收记录数差异非常刺眼Sink[0]的 numRecordsIn 在单位时间内大约是 8000 条而Sink[1]到Sink[7]的 numRecordsIn 几乎为零。这说明并行度是配置上了但数据并没有均匀分发给 8 个写入通道。并行度“生效”和并行度“被用起来”是两件事。Flink 里 Sink 接收的数据路由取决于上游输出流的 channel 分配如果上游只有 1 个并行实例在发送数据那下游即便有 8 个 subtask也只会有一个或极少数subtask 能收到数据。2.2 第二步反压监控暴露出的关键链路接着看 Web UI 的 Backpressure 监控状态如下Source 算子HIGH它在源源不断输出但下游处理不过来时产生的背压信号又传导了回来Sink 算子OK因为真正干活的只有 1 个 subtask对它来说负载并不算高中间没有任何算子Source 直连 Sink。这意味着反压信号主要发生在 Source 与 Sink 之间。我又翻了一下 Runtime Metrics看 Source 算子的numRecordsOutPerSecond始终在 8000 上下Sink 算子的numRecordsInPerSecond总和也是 8000 上下。到这里我基本能确定一个结论整条链路的数据吞吐上限就是 8000 行/s这个上限不是 Sink 决定的而是 Source 决定了的。2.3 第三步用一次临时实验直接验证为了让证据链更完整我在 Source 和 Sink 之间临时插入了一个rebalance()算子。source.rebalance().addSink(dorisSink).setParallelism(8);rebalance()会强制用轮询方式把数据均匀分配到下游每个 subtask。如果瓶颈真的在 Sink 连接池或写入方式上这一步之后吞吐应该有明显提升。结果依然稳定在 8000 行/s只是 UI 上看到 8 个 Sink subtask 都有数据进来了但总吞吐没变。这个实验结果给出了两个信息Sink 端具备处理更高吞吐的能力数据能平均分到 8 个 subtask但它们每秒总共处理的还是那 8000 行Source 端每秒只能吐 8000 行新的瓶颈从“sink subtask 分布不均”转移成了“source 每秒产出总量固定”。到此我可以确定问题本质上是一个结构性限制FlinkCDC 的 MySQL Source 在当前配置下只能提供一个并行实例的数据流。后来我翻源码验证了这一点也发现这是一种普遍存在的架构约束而不是一个能靠参数解决的普通 Bug。3. 根因拆解并行度配置背后的三重锁既然定位到了 Source 端就需要把“为什么 Source 只能单通道产出”这件事彻底讲清楚。很多读者看到这里可能会疑惑FlinkCDC 不是有增量快照框架支持多并行度快照吗为什么还是会被锁死下面拆成三点来说。3.1 Source 端的“单活动 split”限制Flink CDC 2.x 在引入增量快照算法Incremental Snapshot之后全量快照阶段确实可以把整张表按照主键切分成多个 chunk每个 chunk 由不同的 subtask 并行读取这一点在存量数据很大时有明显效果。但增量阶段是另一套逻辑表的 binlog 是一个顺序的、全局有序的变更流为了保证事务顺序和数据一致性在任意时刻只能有一个 split 负责订阅并解析 binlog。也就是说即使你的任务在快照阶段跑了多个并行度一旦进入增量阶段整个 Source 算子实际处于“单活动 split”工作模式。Flink UI 上 Source 并行度可能显示为 4 或 8但真正干活的只有 1 个流。数据源是单水龙头下游装再多水龙头总流量也只能等于单水龙头的出水能力。3.2 “Sink 并行度 写入并发度”的认知误区很多人遇到同步性能问题时第一反应就是调大 Sink 并行度这是一张安全牌但也是一张经常无效的牌。setParallelism(8)的含义是 Sink 算子会创建 8 个并行的 subtask 实例每个 subtask 会尝试与目标库建立连接并执行写入。但每个 subtask 能拿到多少数据取决于上游数据流的分区情况而不是取决于 subtask 的数量。如果上游只有一个并行实例在发送数据下游 8 个 subtask 里只会有一个被持续填充数据其余都处于空转状态。你可以把整条链路想象成一条单车道高速路终点有 8 个收费站。无论你把收费站从 1 个增加到 8 个单位时间内能到达终点的车流量仍然取决于那条单车道的通行能力。真正能解决问题的办法是拓宽车道或者让车辆先汇聚到中间的大型转运中心再分多条路去往收费站。3.3 参数链路上的“无效优化”陷阱在 Source 单通道锁定的情况下很多常见优化参数不会产生质变甚至会产生误导sink.buffer-flush.max-rows调大后Sink 只是把同样数量的数据攒成大包再写但数据总量没变省下的只是批提交开销sink.buffer-flush.interval拉长后反而会让延迟变高吞吐并没有明显提升增大 JDBC 连接池或 Doris 的 buffer size解决的是并发连接争抢问题在只有 1 个 subtask 在写入时基本没有帮助。所以当你发现调整这些参数都无效时不要继续在这个维度上死磕大概率方向已经错了。3.4 对比反例为什么 Kafka Source 很少遇到这个限制理解这个问题最好的方式是找个反面对比。如果你从 Kafka 读取数据Source 并行度可以设置为 Topic 分区数16 个分区就可以让 16 个 subtask 同时消费。此时数据流天然是分区的、并行的下游 Sink 只要和分区数匹配就能利用上多个写入通道。MySQL CDC 和 Kafka 最大的差异就在这里Kafka 本身是分布式消息队列分区是物理存在的binlog 则是单一文件流所有并行读取最终都要归约到这个单一顺序流上。这是架构层面决定的所以靠调参数很难绕过必须从架构设计上寻找突破口。4. 解决方案读写解耦与分区器组合拳根因清楚了下面就是方案选型。我这里按推荐优先级给出三种实际验证过的解决路径。4.1 方案一读写解耦拆成两个 Job 串接 Kafka这是我在生产环境实测效果最好的方案也是目前大规模 CDC 同步场景里最通用的架构。改之前的链路MySQL binlog - FlinkCDC Source(1) - Doris Sink(8)改之后的链路Job1: MySQL binlog - FlinkCDC Source(1) - Kafka Sink(8) Job2: Kafka Source(8) - Doris Sink(8)核心思路让 Source 端和 Sink 端不再直连中间隔一层 Kafka。Kafka 天然适合做削峰填谷和并行数据分发上游多少个 Kafka Sink 实例可以自由设置下游 Kafka Source 的并行度可以设置为 Topic 分区数从而让多个 subtask 同时消费并写入目标库。操作步骤第一步创建 Kafka Topic。分区数是整个方案成败的关键我推荐按照“下游预期并行度 × 1.5 到 2”来规划。下游 Job2 的并行度是 8Topic 分区数设置为 16。kafka-topics.sh --create --bootstrap-server kafka-1:9092 \ --topic ods_order_cdc \ --partitions 16 \ --replication-factor 3分区数不能设得太小否则下游并行度会被限制在分区数以内也不能设得过大否则 Kafka 自身因为文件句柄、副本同步带来的开销会变得明显延迟反而上升。第二步编写 Job1。Job1 只负责两件事读 MySQL CDC写 Kafka。这个 Job 的 Sink 并行度建议设置在 8小于等于 Topic 分区数。DataStreamString source env.addSource( MySqlSource.Stringbuilder() .hostname(mysql-primary) .port(3306) .databaseList(shop) .tableList(shop.t_order) .username(cdc_user) .password(******) .deserializer(new JsonDebeziumDeserializationSchema()) .startupOptions(StartupOptions.initial()) .serverTimeZone(Asia/Shanghai) .build() ); source.addSink(KafkaSink.Stringbuilder() .setBootstrapServers(kafka-1:9092,kafka-2:9092,kafka-3:9092) .setRecordSerializer(KafkaRecordSerializationSchema.builder() .setTopic(ods_order_cdc) .setValueSerializationSchema(new SimpleStringSchema()) .build()) .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE) .setTransactionalIdPrefix(cdc-to-kafka) .build()) .setParallelism(8);注意Kafka Sink 的并行度不要超过 Topic 分区数否则多出来的 subtask 会发现没有可写的分区白白浪费资源。第三步编写 Job2。Job2 从 Kafka 读取写入 Doris。Kafka Source 的并行度和 Topic 分区数对齐这里设置 8Doris Sink 也设置 8。DataStreamString stream env.fromSource( KafkaSource.Stringbuilder() .setBootstrapServers(kafka-1:9092,kafka-2:9092,kafka-3:9092) .setTopics(ods_order_cdc) .setGroupId(doris-sink-group) .setStartingOffsets(OffsetsInitializer.committedOffsets()) .setValueOnlyDeserializer(new SimpleStringSchema()) .build(), WatermarkStrategy.noWatermarks(), kafka-source ).setParallelism(8); stream.addSink(dorisSink) .setParallelism(8);这套方案的收益是最直接的因为下游 Job2 有 8 个并行消费线程每个线程写入 Doris 时都不再受上游单通道限制。我实测改造后吞吐从 8000 行/s 提升到 62000 行/s 左右提升接近 8 倍。这个方案还有一个附加好处两个 Job 可以独立扩容。如果某天 Doris 写入成为瓶颈只需要调大 Job2 的并行度并对着 Kafka Topic 增加分区数不需要重新拉起 CDC 读取任务全程不用停线上同步。4.2 方案二自定义分区器解决 Sink 端数据倾斜如果你不想引入 Kafka但你的 Source 本身已经是多并行度例如从 Kafka 读、或者使用支持多并行度快照的版本那 Sink 端吞吐上不去的主要矛盾可能是数据倾斜。一个典型的错误写法是按照业务热点字段 keyBy。比如同步订单数据时有些用户大客户的单量特别大如果你按userId做 keyBy那么所有大客户的变更记录都会集中到同一个 subtask造成严重的倾斜。CDC 场景下最稳妥的做法是按主键的哈希值做分区保证同一条记录永远落到同一个分区同时尽量让不同主键均匀分布。如果对顺序性要求高不能使用随机散列因为同一主键的 update/delete 乱序到达目标库会导致最终数据不一致。source .keyBy(new KeySelectorString, String() { Override public String getKey(String value) throws Exception { // 假设解析出主键 orderId按键的 hash 均匀分布 return orderId; } }) .addSink(dorisSink) .setParallelism(8);如果你对同一条记录的变更顺序并不敏感比如目标表只保留最终状态也可以考虑在 Source 和 Sink 之间插入rebalance()用轮询强制重新分区。但这在实时同步任务里要慎用因为一旦出现乱序纠错成本非常高。4.3 方案三不拆 Job最大化单 Sink 写入能力如果业务量很小或者实在不想为了一个同步任务引入新的 Kafka 集群可以在单 Job 内做参数层面的优化。这个方案的适用上限比较低通常只能提升 10%-30%但胜在改动小。在 Doris Connector 中可以重点调这几个参数调大 stream load 的 batch 大小减少提交次数调大 buffer-size降低小包高频写入的开销调大sink.max-retries的同时关注目标库侧是否出现导入版本冲突将 checkpoint 间隔从默认值调整到 30 秒以上给两阶段提交留足时间。需要注意这个方案无法突破 Source 单通道的总量瓶颈。如果 Source 端每秒只能吐 8000 行Sink 再好也是巧妇难为无米之炊。4.4 三个方案的适用场景对比方案架构变化预期吞吐提升实施成本适用场景读写解耦串接 Kafka引入 Kafka 中间层5-10 倍中数据量大、生产环境、长期运行自定义分区器无架构变化视倾斜程度 30%-300%低Source 多并行度但存在热点数据Sink 参数优化无10%-30%低小数据量、临时任务我个人的建议是只要业务量会持续增长就直接上方案一。这套架构不仅仅提升了单一的同步性能后面你要做多表汇聚、多目标分发Kafka 这个中间层都会变成基础设施早晚都要建。5. 实测效果与参数调优清单方案一落地之后我记录了改造前后的具体数据也沉淀了一份可以直接抄的参数清单下面完整列出来。5.1 改造前后的性能对比指标改造前改造后稳定吞吐约 8000 行/s约 62000 行/s峰值吞吐8500 行/s75000 行/s全量 800 万行耗时约 18 分钟约 3 分钟数据同步延迟分钟级秒级TaskManager CPU 利用单核打满其余空闲8 核基本均衡为什么是接近 8 倍而不是严格的 8 倍因为 Kafka 消费端每次拉取是一批数据Doris 的 stream load 每次提交也有批大小限制这些批处理环节会有一定吞吐损耗但整体上接近线性扩容。5.2 可直接抄取的参数清单配置项参数值说明Kafka Topic 分区数16下游并行度的 1.5-2 倍Job1 Source 并行度1受 binlog 单流限制无法突破Job1 Kafka Sink 并行度8不要超过 Kafka Topic 分区数Job2 Kafka Source 并行度8对齐 Kafka Topic 分区数Job2 Doris Sink 并行度8与 Kafka Source 保持一致checkpoint interval30s过长延迟变高过短事务开销大Kafka 事务超时900000ms对应 transaction.max.timeout.msDorisSink batch 大小1000 行以上减少 stream load 提交频率需要提醒一点Kafka 事务相关的参数和 checkpoint 间隔是联动的。KafkaSink使用 EXACTLY_ONCE 时每个 checkpoint 周期会开启一个事务transaction.timeout.ms必须大于 checkpoint 间隔否则写入会直接报错。生产环境建议把 Kafka 的transaction.max.timeout.ms调大不然频繁的大事务会被 broker 拦截。5.3 怎么验证瓶颈真的被解除了改造完成后不要只看吞吐数字还需要做三个检查第一看 Flink UI 反压状态。改造后 Job1 和 Job2 都应该处于 OK 状态Source 和 Sink 之间不再有强烈背压。第二看 Kafka 消费延迟。执行命令kafka-consumer-groups.sh \ --bootstrap-server kafka-1:9092 \ --describe \ --group doris-sink-group正常情况下 LAG 应该稳定在一个很小的值附近如果 LAG 持续增长说明 Job2 的消费速度跟不上 Kafka 的写入速度还需要继续扩大 Job2 并行度。第三看 Doris 侧的实际导入监控。观察各 BE 节点的 tablet 写入流量正常情况下各节点的写入负载应该是接近均衡的。如果某个节点独高说明数据分区策略还需要优化。6. 复盘与后续避坑问题解决之后我又回看了整次排查过程有几个很容易被忽略的点值得单独拎出来说。第一个经验是遇到“提高并行度但性能不变”这类问题先看 Flink Web UI 的numRecordsInPerSecond和numRecordsOutPerSecond这两个指标不要直接改配置。这两个指标能最快告诉你瓶颈到底在哪个算子。如果挂了一晚上吞吐都没变源头大概率已经被锁死了。第二个经验是Source 单通道限制是架构问题不要在上面浪费太多时间调参数。我后来查了不少社区讨论发现很多人也遇到同样的问题在server-id、fetchSize、debezium参数上反复调结果收效甚微。正确的做法不是优化 Source 的读取速度而是把数据通道从“单车道”改成“多条车道并行”Kafka 中间层就是为了解决这个问题而存在的。第三个经验是拆成两个任务之后Kafka 的分区策略一定要设计好。最简单可靠的策略是直接按主键的哈希做路由这样同一行的变更会始终进入同一个分区目标库端不会因为重合记录产生版本错乱。如果你用轮询分区很可能在异常恢复或重复消费时出现乱序导致目标库数据短暂不一致。第四个经验是Kafka Sink 的 EXACTLY_ONCE 语义是有代价的。开启事务性写入之后每个 checkpoint 周期都要开启一个 Kafka 事务当分区数很多且 checkpoint 间隔很短时broker 端的事务协调会消耗不少资源。我实测下来如果业务对“精确一次”要求没那么严格比如同步到数仓做后续离线修正使用 AT_LEAST_ONCE 可以把吞吐再往上提一截。至于重复数据在目标库通过主键去重可以兜底。最后分享一个我自己踩过的坑。有一次 Job2 的 Kafka Source 并行度被同事改成了 20但 Topic 分区数只有 16。结果启动后一直有 4 个 subtask 处于空闲状态另外 16 个 subtask 忙得不行吞吐也没有提升。这个细节在界面上看不直观因为并行度确实显示 20/20但消费者数量大于分区数时多余的 subtask 永远等不到数据。所以记得随时检查“并行度、分区数、实际负载”三者是否匹配。这个 Bug 解决之后我又用同样的读写解耦方案处理了好几套同步链路包括订单库、库存库和用户行为日志库。基本套路都是“MySQL CDC 进 Kafka再由下游任务自由消费、自由扩容”再也没有被 Source 单并行度锁死过。如果你现在也被同步性能卡得难受建议先按文中的排查链路走一遍确认瓶颈之后再考虑是否引入 Kafka 这一层。
返回列表