ARTICLE DETAIL

资讯详情

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

Flink CDC同步性能瓶颈:Sink并行度未生效的排查与调优

Flink CDC同步性能瓶颈:Sink并行度未生效的排查与调优 先说结论这个坑我踩了整整一天如果你也遇到FlinkCDC同步任务吞吐上不去、延迟持续增长、并行度怎么调都没用的情况大概率和我遇到的是同一个原因——写入端的并行度并没有真正生效。先说下我当时的场景源端是MySQL通过FlinkCDC版本用的2.3.0实时解析binlog同步到下游目标库。整体链路是MySQL - FlinkCDC - Transformer - Sink用的是DataStream API开发的跑在Flink on YARN上作业Manager资源给得不算低。业务量上来之后发现同步任务开始积压延迟从几十秒一路涨到十几分钟。当时第一反应是source读binlog不够快于是把source并行度从1调到4结果吞吐纹丝不动。后来又调全局并行度、加TaskManagerCPU和内存明明都还有余量可吞吐就是卡在两万条每秒上不去了。后来我把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 4sink节点可能还是会变成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修改后的代码片段DataStreamOrderRow 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子任务闲置假如你的上游并行度是4Sink并行度也设成了4但每个Sink SubTask的写入量差距特别大性能依然上不去。这是因为Flink默认的Forward模式下上游每个Subtask固定把数据发给下游同一个Subtask并不会自动负载均衡。解决的方案是在Sink之前加一个rebalance()强制数据轮询分配到下游每个SubtaskparsedStream .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以内。以下数据是稳定运行后的平均表现阶段配置平均吞吐条/秒同步延迟备注调优前全局并行度4Sink未显式设置约2.1万持续积压最高达15分钟Sink实际并行度为1第一步调优Source 4Sink 4setParallelism(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未显式设置且默认并行度为1keyBy断链给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这样监控面板上一眼就能看出并行度是不是已经生效。多花几分钟做这些细节比事后排查一天要划算得多。
返回列表