先说结论:如果你的 FlinkCDC 数据同步任务遇到同步性能无法提升、怎么调 Sink 并行度吞吐都纹丝不动的情况,大概率问题不在 Sink 端,而是整条链路的写入并行度被上游 Source 的单通道给锁死了。这个坑我踩了一整天才彻底定位,当时从发现写入 Doris 的吞吐上不去,到最后通过读写解耦方案把性能提升了接近 8 倍,中间那一整条排查链路,我认为很值得单独写一篇复盘。不管你现在用的是 FlinkCDC 2.x 还是 3.x,只要架构还是“MySQL CDC 直连目标库”,这篇文章的内容就都适用,其中包含可照抄的参数配置、反压判定技巧,还有几个特别容易误判的细节。
1. 现象还原:Sink 并行度调到 8,吞吐反而没变化
先交代一下任务背景。我当时要同步的是线上一个订单库,MySQL 单实例,存量订单表接近 800 万行,日均新增 50 万左右。目标库是 Doris,用来做实时 OLAP 报表。Flink 版本 1.16.2,FlinkCDC 用的 2.3.0,代码基于 DataStream API,没有走 SQL Client,这样方便后面控制并行度和分区逻辑。
任务的初始拓扑非常简单:
MySQL binlog -> FlinkCDC Source -> Doris SinkSource 端按照官方推荐的方式配置,用MySqlSource的 Builder 构建,启动模式用的initial(),因为第一次要同时处理全量存量数据和增量日志。Sink 端用的是 Doris Connector,为了拉高写入并行度,我显式指定了.setParallelism(8)。
DataStream<String> source = env.addSource( MySqlSource.<String>builder() .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 | 反压 High |
| 4 | 默认 100 | 约 7600 行/s | 反压 High |
| 8 | 默认 100 | 约 7900 行/s | 反压 High |
| 8 | 5000 | 约 8100 行/s | 无实质改善 |
| 8 | 2000 | 约 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 的并行度是 8,Topic 分区数设置为 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 分区数。
DataStream<String> source = env.addSource( MySqlSource.<String>builder() .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.<String>builder() .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 分区数对齐,这里设置 8,Doris Sink 也设置 8。
DataStream<String> stream = env.fromSource( KafkaSource.<String>builder() .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 KeySelector<String, 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 行/s | 75000 行/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 interval | 30s | 过长延迟变高,过短事务开销大 |
| Kafka 事务超时 | 900000ms | 对应 transaction.max.timeout.ms |
| DorisSink 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 这一层。