先说结论:这个坑我踩了整整一天,如果你也遇到FlinkCDC同步任务吞吐上不去、延迟持续增长、并行度怎么调都没用的情况,大概率和我遇到的是同一个原因——写入端的并行度并没有真正生效。
先说下我当时的场景:源端是MySQL,通过FlinkCDC(版本用的2.3.0)实时解析binlog,同步到下游目标库。整体链路是MySQL -> FlinkCDC -> Transformer -> Sink,用的是DataStream API开发的,跑在Flink on YARN上,作业Manager资源给得不算低。业务量上来之后,发现同步任务开始积压,延迟从几十秒一路涨到十几分钟。当时第一反应是source读binlog不够快,于是把source并行度从1调到4,结果吞吐纹丝不动。后来又调全局并行度、加TaskManager,CPU和内存明明都还有余量,可吞吐就是卡在两万条每秒上不去了。
后来我把Flink Web UI打开,逐个算子看Subtask数量,才彻底意识到问题:并行度根本没有传递到Sink端。所有数据最终都挤在一个Sink SubTask里,单连接单线程写目标库,性能上限就是一条单链路的写入能力。这篇文章就围绕这个“并行度丢失”问题,把排查过程、根因分析、解决方案和验证数据全部梳理一遍。不管是DataStream API还是Flink SQL,核心思路是相通的,看完你应该能直接照着排查。
1. 问题表象与初步排查路径
1.1 同步延迟飙升但资源利用率很低
我这里跑的是一张订单大表,日增量大概是千万级别。出问题之前,同步任务延迟基本维持在一两分钟以内,某天业务大促之后明显感觉不对劲,打开监控面板发现同步延迟曲线直接掉头向上,而整条链路的CPU使用率不到30%,网卡流量也不高。
这种“资源没用满但性能上不去”的现象,是最让人头疼的。直觉告诉我一定有个隐性的单点瓶颈,而不是资源不够。于是我先做了两件事:第一,确认source的binlog读取是否有积压,如果binlog消费正常说明源头不缺数据;第二,在Flink Web UI上查看各个算子的BackPressure状态,看数据到底堵在哪一环。
第一轮排查下来,source端的Records Sent Count和binlog位点推进都很正常,说明MySQL binlog读取端没有问题。但是Sink端的收数据速率明显低于source,而且最关键的是——Sink算子只有一个Subtask。数据在source端已经被拆到4个并行子任务里处理了,到了Sink端却被强制收敛成一个线程写入,这个单线程瞬间就成了整条链路的瓶颈。
1.2 确认瓶颈算子的方法
确认瓶颈位置有几个比较实用的办法,不用瞎猜:
- 看反压状态:Flink Web UI的BackPressure选项卡,如果某个算子显示
HIGH,说明数据在它上游堆积了。但注意,这里反压表示的是上游算子因为下游处理不过来而变忙,所以看到HIGH的算子,瓶颈往往在它的下游。 - 看Busy%比例:进入某个算子的subtask详情,如果
busy%长期在90%以上,基本可以断定这个算子就是性能瓶颈。 - 看记录堆积指标:
numRecordsIn和numRecordsOut的差值,如果某个算子的输入远大于输出,它或者它的下游必然存在问题。
实测的结果是:Source端4个Subtask,每个都在拼命输出;而Sink端只有1个Subtask,它的numRecordsIn持续上涨、numRecordsOut却上不去——典型的sink消化能力不足,而且这个sink根本没有并行执行。
注意:Flink Web UI上显示的
Number of Sub-Tasks这一列,是最直观的判断依据。如果你发现某个算子Subtask数量一直为1,且代码里又没显式给它设置并行度,那就要警惕了,并行度可能在算子链路中被“稀释”了。
2. 根因剖析:并行度在Flink算子链路中的传播规则
2.1 并行度设置的三个层级
Flink里的并行度可以设置在很多地方,默认生效的优先级是这样的(从上到下优先级逐渐降低):
- 算子级别:
map.map(…).setParallelism(4),最精确,只对当前算子生效。 - 执行环境级别:
env.setParallelism(4),对整个作业里所有没有显式设置并行度的算子生效。 - 提交参数级别:
-p 4,提交作业时通过命令行参数指定。 - 配置文件级别:
flink-conf.yaml里的parallelism.default,兜底值。
这里有个关键点:如果你在算子A后面不设置任何并行度,Flink不一定把算子A的并行度传给算子B。在DataStream API里,两个相邻算子默认会尝试合并成一个算子链(operator chain),合并后它们共享同一个Subtask里的线程,并且使用同一个并行度。但如果中间的算子有keyBy、rebalance、rescale或者startNewChain之类的操作,链就会被断开,下游算子的并行度会回落到执行环境的默认并行度,如果你没设置过全局并行度,那就默认是1。
在我这个场景里,Source端设置了并行度4,但Source算子后面接了filter和map,其中map里做了一次keyBy操作——为了按主键分组处理数据,我当时以为这样能保证数据不乱序。结果就是:keyBy把operator chain断开了,而map之后的Sink没有显式设置并行度,Flink直接给Sink分配了默认并行度1。所有数据在keyBy之后,都被路由到同一个Sink SubTask里,单线程往下写。
2.2 SQL作业里的并行度传播差异
如果你是写Flink SQL做同步,情况会有些不同。Flink SQL天然会把整个作业拆成Source、Transform、Sink几个逻辑节点,每个节点的并行度既可以单独配置,也可以继承作业全局并行度。比如你在SQL Client或者TableAPI里这样设置:
-- 全局并行度 SET 'parallelism.default' = '4'; -- 或者单独指定某个sink的并行度 SET 'sql.sink.parallelism' = '4';但有个点坑过很多人:Flink SQL里设置了parallelism.default = 4,sink节点可能还是会变成1,原因要看具体连接器实现。比如某些JDBC Connector的Sink,内部需要维护数据库连接池,设计上默认就是单并行度写,避免多线程并发写导致连接数暴涨或者主键冲突。你必须去查看对应连接器的官方文档,确认它是否支持并行写入,以及并行度参数的准确名称。
以下几种场景下,并行度特别容易被重置为1:
- 连接器官方默认不支持并行写(比如部分版本的JDBC sink、HBase sink)。
- 代码里对Sink调用了
.setParallelism(1)但自己忘了。 - SQL作业里没有显式配置sink并行度,且连接器内部强制单线程。
- 资源不足,TaskManager总Slot数小于你设置的并行度,Flink不得不按资源上限分配。这种情况不会报错,但在Web UI上看起来就是并行度没生效。
2.3 为什么单并行度Sink会造成同步性能天花板
单并行度意味着整个写入链路只有一个线程,这个线程要完成序列化、建连、发送、等待目标库确认等所有工作。即使你的批处理size设置得很大,单线程的网络往返延时和数据库锁等待也会成为硬约束。
我做了个粗算:假设单次网络往返延迟是2ms,一个批次写1000条,那么每秒理论上的批次上限是500次,也就是50万条。听起来不低对吧?但如果你用的小批次(比如每50条提交一次),单线程每秒最多只能提交200次,吞吐直接掉到1万条每秒。我当时的Sink配置正好是每批200条就提交一次,单线程吞吐被锁死在2万条/秒附近,和监控面板上的数据完全吻合。
所以,凡是做实时同步,一定要把写入并行度当成一等公民来设计,否则哪怕source能读的再快,最终还是会卡在写入单点上。
3. 问题定位:Flink Web UI实操看并行度分配
3.1 从Overview页面和Job Graph里找异常
我实际排查的时候,Flink Web UI帮了大忙。点开作业详情,在Job Graph页面可以看到每个算子的Subtask个数。我当时一眼就发现,Source: MySQL CDC下面有4个小格子,但Sink: JdbcSink下面只有1个小格子——并行度不对称一目了然。
然后点进Sink算子,进Subtask列表页面,看Bytes Received和Records Received。如果一个Sink Subtask同时接收来自上游4个不同Subtask的数据,那它的Records Received数值会是上游总数的总和,这更加印证了数据在最后一步被强制归并了。
接着我打开Metrics页面,查看Sink算子的numRecordsOutPerSecond。注意只看瞬时值,波动比较大的多采样几轮。如果每秒钟输出稳定在某个值附近上不去,比如2万左右,那基本可以断定写入端性能上限就在这了。
3.2 用反压和Watermark辅助验证
除了直接看Subtask数量,还可以用反压状态来交叉验证。当时Sink算子上游几个算子全部显示HIGH反压,而Sink本身的busy%高达97%,说明这个单线程Sink已经满负荷运转了。水印(Watermark)也在Sink算子上有严重的滞后,意味着事件时间线被阻塞住了。
水印滞后的计算方式很直观:如果Source端已经解析到binlog里的最新位点,但Sink端水印还停留在十分钟前,那中间的差就是目前同步的延迟时长。这个指标配合反压状态,基本可以拼出完整的故障图。
还有一个容易忽略的点:如果你用了
keyBy做分区,那即使Sink并行度不为1,也很可能出现数据倾斜——某些key的数据量大,对应的Subtask忙死,其他Subtask闲得没事干。我当时Sink并行度为1,所以不存在数据倾斜问题,但如果你调大并行度之后性能还是没变,一定要回过来检查这一项。
4. 解决方案:三步真正提升写入并行度
定位到问题之后,解决起来反而简单了。我的核心思路有三个:让Sink算子拥有真正的并行度、让数据在多个Sink SubTask之间均匀分布、让每个Sink SubTask的写入性能最大化。
4.1 显式给Sink算子设置并行度
最直接的改法,就是在代码里给Sink算子显式设置并行度。我这里用的是DataStream API,修改后的代码片段:
DataStream<OrderRow> parsedStream = sourceStream .filter(new OrderFilter()) .keyBy(order -> order.getOrderId()) .map(new OrderTransform()); // 关键:显式指定sink并行度,不再继承默认值 parsedStream.addSink(new OrderJdbcSink()) .setParallelism(4);这里有两个容易踩的坑:
- 如果你
keyBy之后的数据,仍然需要保证同一个主键能稳定落到同一个Sink SubTask里,那么keyBy和setParallelism(4)要配合使用,Flink默认会按key的hash值取模分配给下游Subtask,相同的key一定会进同一个Subtask,这一点可以放心。但如果你不keyBy直接sink.setParallelism(4),数据大概率会走Forward模式,只有一个上游Subtask会连到特定的下游Subtask,仍然可能造成部分Sink SubTask空闲。 - 如果下游是目标数据库,而且目标库对单连接写入有瓶颈,那你把Sink并行度调大后,每个Subtask都会建立独立的数据库连接。这时候要确认目标库的最大连接数够不够,连接池配置是否合理。连接数不够的话,并行度调了反而会报连接超时。
4.2 对数据流显式重分区,防止Sink子任务闲置
假如你的上游并行度是4,Sink并行度也设成了4,但每个Sink SubTask的写入量差距特别大,性能依然上不去。这是因为Flink默认的Forward模式下,上游每个Subtask固定把数据发给下游同一个Subtask,并不会自动负载均衡。
解决的方案是在Sink之前加一个rebalance(),强制数据轮询分配到下游每个Subtask:
parsedStream .rebalance() .addSink(new OrderJdbcSink()) .setParallelism(4);rebalance()底层是通过Round-Robin方式把数据均匀分发到下游所有并行实例。代价是有一定的网络序列化开销,但换来的负载均衡非常值得。
如果你的业务场景对数据顺序有强要求,那么rebalance()要慎重。多并行度Sink本质上会把写入顺序打散,如果下游没有主键约束或者去重逻辑,很容易出现乱序写入。建议在目标表上设计好主键,并且尽量使用幂等写入模式,也就是基于主键做upsert,这样即使乱序到达,最终结果也是正确的。
4.3 Sink端写入参数调优:批次大小、刷新间隔、连接池
并行度解决了“多线程写”的问题,但每个线程自身的写入效率也值得调优。这里以JDBC Sink为例,几个关键参数:
- 批次大小(batch size):我调优前是200条提交一次,调大到了1000条。提交频率降下来,事务开销明显减少。但注意批次不能无限大,太大会导致单个事务执行时间过长,数据库锁持有时间变长,反而降低吞吐。
- 刷新间隔(flush interval):如果你希望延迟不能太高,可以保持一个较低的刷新间隔,比如3秒强制刷一次。实时同步场景需要在低延迟和高吞吐之间做个平衡,我最后用的是“batch size达到1000或时间达到3秒,谁先到谁触发”。
- 连接池大小:如果你在Sink里用了连接池,并发度提升后连接池最大连接数也要相应提升。我用的Druid连接池把
maxActive从10调到了50,否则4个Sink SubTask加其他任务抢连接,很容易把连接池耗尽。
代码里配置连接池和提交频率时,建议把参数统一放到配置文件里,方便不同环境调整。我当时的做法是:
JdbcExecutionOptions execOptions = JdbcExecutionOptions.builder() .withBatchSize(1000) .withBatchInterval(Duration.ofSeconds(3)) .build(); JdbcConnectionOptions connOptions = JdbcConnectionOptions.builder() .withUrl(jdbcUrl) .withDriverName("com.mysql.cj.jdbc.Driver") .withUsername(username) .withPassword(password) .build(); sink = JdbcSink.sink( insertSql, new OrderJdbcMapper(), execOptions, connOptions );4.4 目标库侧配合:索引、锁等待、写入模式
并行度调上去了,目标库可能成为新的瓶颈。这时候不能只盯着Flink,要对目标库做几个常规检查:
- 主键和索引:目标表的索引越少写入越快,尤其要避免在写入表上有多个二级索引。每条insert都会维护所有索引,二级索引多会拖慢写入速度。可以考虑在同步期间暂时禁用非必要索引,数据补完再重建。
- 锁等待:如果目标表同时有业务查询在跑,写入可能经常被行锁或间隙锁阻塞。建议给同步任务使用的数据库账号单独设置合理的锁等待超时时间,避免长时间卡死。
- 写入模式:JDBC默认是INSERT,如果源库有数据更新,建议改成
INSERT ... ON DUPLICATE KEY UPDATE或者REPLACE INTO,这样多个并行Sink SubTask即使写同一个主键,也不会因为重复报错而失败。
我当时对目标表做了一次重建二级索引的操作,光这一步,单SubTask写入速度就提升了将近30%。这个优化经常被忽略,但其实性价比很高。
5. 调优后的效果对比与参数记录
5.1 调优前后实测数据
我以订单数据同步为例,连续跑了3个小时,测试环境是4个TaskManager,每个TaskManager 2个Slot,总共8个Slot。目标库是单独的MySQL实例,网络延迟在1ms以内。以下数据是稳定运行后的平均表现:
| 阶段 | 配置 | 平均吞吐(条/秒) | 同步延迟 | 备注 |
|---|---|---|---|---|
| 调优前 | 全局并行度4,Sink未显式设置 | 约2.1万 | 持续积压,最高达15分钟 | Sink实际并行度为1 |
| 第一步调优 | Source 4,Sink 4(setParallelism(4)) | 约4.6万 | 降至2~3分钟 | 仍有轻微积压 |
| 第二步调优 | 增加rebalance() 均匀分发 | 约5.2万 | 基本稳定在1分钟以内 | 数据无明显倾斜 |
| 第三步调优 | 批次200->1000,连接池增大,目标库索引优化 | 约8.5万 | 稳定在30秒以内 | 接近目标库单实例写入上限 |
可以看到,每步调优都在叠加效果,但后面几个阶段性能提升也在放缓,最终8.5万条/秒接近了这台目标库单实例的写入天花板。如果再往上提,要么做目标库分库分表,要么换写入性能更强的存储引擎,单纯调Flink参数已经没有太大收益了。
5.2 稳定性验证与回滚预案
调优之后不能光看吞吐涨了多少,还必须观察是否稳定。我连续观察了三个高峰期,确认以下指标全部正常:
- TaskManager的CPU使用率稳定在60%左右,没有持续飙高。
- GC时间没有异常增长,Full GC次数保持低位。
- 目标库的线程数、连接数都在合理范围内,没有锁等待超时。
- 同步延迟在高峰过后能迅速回落到30秒以内,具备自愈能力。
另外我还留了一手回滚方案:把所有调优参数都写在了配置中心,通过动态配置可以随时切回旧值。万一新参数在高负载下出现意外,我可以快速回退,而不需要重新提交Flink作业。这套配置化的方式强烈建议你也采用。
6. 常见问题速查与避坑建议
6.1 几个容易踩的坑
我把自己和身边同事踩过的坑整理成了一个表,方便你对照排查:
| 问题现象 | 可能原因 | 解决方案 |
|---|---|---|
| Sink并行度始终为1 | 未显式设置,且默认并行度为1;keyBy断链 | 给Sink显式.setParallelism(n) |
| 设置了并行度但吞吐不变 | 数据倾斜,部分Subtask空闲 | 加rebalance()或rescale()均匀分发 |
| 并行度调大后目标库连接数爆掉 | 连接池和数据库最大连接数没协调好 | 增大连接池上限,同时检查数据库max_connections |
| 并行写后出现主键冲突 | 多个Sink线程同时写同一主键,且非幂等模式 | 换成upsert写入模式,目标表加主键 |
| 调大批次后延迟升高 | 批次太大,flush间隔过长 | 设batch size和时间间隔双重触发条件 |
| 并行度明明设置了,重启后失效 | 作业提交时被-p参数或配置文件覆盖 | 确认提交参数、Flink配置和代码设置三者一致 |
6.2 建议保留的黄金排查路径
以后再遇到类似的同步性能问题,我的建议是按下述路径走,可以少走很多弯路:
- 先看Flink Web UI,逐算子核并行度,把每个算子的Subtask数量记下来。
- 再点进瓶颈算子看反压状态和
busy%,判断性能瓶颈是在当前算子还是下游。 - 用监控面板看整个链路的吞吐、延迟、GC、连接数,排除外部系统瓶颈。
- 确定是Sink并行度问题后,先小范围调整(比如2并行度),验证有效再加到4、8。
- 每次调优只改一个变量,不要同时改并行度、批次、索引,否则你根本不知道是哪个改动起了作用。
逐步验证是我这次排查中最受益的习惯。如果你同时调了三四个参数,性能涨了,但根本说不清是哪个参数起的关键作用,下次遇到问题还得从头试。
写在最后的一点私货
FlinkCDC这种场景,“实时同步”看起来是连接器在干活,实际上瓶颈大概率在接入端和写出端。这次排查最大的收获是让我形成了一个习惯:任何Flink作业提交前,先在Web UI上把算子的并行度完整过一遍。并行度丢失这种问题,运行时表现隐蔽,定位起来费时间,但只要确认好每个算子的Subtask数,基本一眼就能识破。
最后再分享一个小技巧:给Sink的每个Subtask命名加一个前缀,比如sink-writer-0、sink-writer-1,这样监控面板上一眼就能看出并行度是不是已经生效。多花几分钟做这些细节,比事后排查一天要划算得多。