免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

以Kafka Connector为线,读懂Flink的checkpoint与Exactly-Once

以Kafka Connector为线,读懂Flink的checkpoint与Exactly-Once Flink源码阅读这个坑我入得比较晚但一入坑我就发现想多了与其漫无目的地翻各种算子和状态机不如找一条主线啃到底。我的选择是Kafka Connector原因很实在——它把Flink的checkpoint、状态恢复、背压、水印传播和Kafka的消费组、分区、事务机制全部串在一起读通它的源码你对整个Flink的运行时机制理解会直接上一个台阶。这篇文章不打算把每个类都念一遍而是挑几条主线讲消费者怎么恢复offset、怎么发现新分区、生产者怎么用两阶段提交实现Exactly-Once以及在真实故障里这些源码知识怎么帮你快速定位问题。1. 为什么源码阅读要选Kafka Connector一份代码两套框架的精华1.1 从一道Flink面试题说起有次我和一个团队做内部分享他们刚面完一个候选人面试官出的题是你的Flink作业消费Kafka的topicKafka那边扩了分区Flink是怎么感知到的新分区从哪个offset开始消费会丢消息吗候选人答得磕磕绊绊最后面试官自己也讲不清细节只能按结论判对错。这道题的答案就藏在FlinkKafkaConsumerBase的代码里。如果你只背结论遇到追问就露馅但你要是读过源码哪怕忘了具体的行号也能把那条链路讲清楚分区发现线程定期调getAllPartitionsForTopics拉取最新分区集合和当前订阅的KafkaTopicPartition集合做对比发现新分区后按启动模式决定起始offset再交给fetcher加入订阅。这条链路里每一环都有代码可以指认面试官一听就知道你是真读过。1.2 读源码前需要建立的Kafka基础知识直接读代码很容易被细节淹没我建议先建立三个坐标。第一分区是调度和容错的最小单位。Kafka的一个分区在Flink看来是一个split一个并行子任务可以消费多个分区但一个分区同一时刻只能被一个子任务消费不能两个子任务抢同一个分区这是语义正确的前提也是后面理解分区发现和offset管理的基础。第二offset的语义要抠准。Kafka里的offset表示下一条待消费消息在分区中的位置Flink checkpoint保存的正是这个下一条的offset不是最后一条已消费的记录。这个细节在FlinkKafkaConsumerBase的注释里写得很明确很多人在调seekToEnd之类的方法时把位置搞混根源就在这。第三Kafka生产者的事务能力是Flink端到端Exactly-Once的底层支柱。Kafka在0.11版本之后支持事务生产者可以把跨分区的写入作为一个原子事务提交Flink基于这个能力实现了自己的两阶段提交。理解这一点你才看得懂FlinkKafkaProducer里那套beginTransaction、preCommit、commit的循环。2. 消费者链路拆解FlinkKafkaConsumer的启动、恢复与数据流转2.1 初始化与状态恢复offset从哪儿来FlinkKafkaConsumerBase这个抽象类实现了CheckpointedFunction和CheckpointListener这是整个消费端和状态生命周期对接的入口。作业启动时initializeState()会通过getRuntimeContext().getUnionListState(descriptor)读取之前checkpoint保存的分区offset快照在较新的版本里使用的是UnionListState老版本里是ListState它们的区别在于并行度变化时状态的重新分发方式。我第一次读这段代码时有个误区一直以为作业重启后Flink会从Kafka的consumer group提交的offset开始消费。实际上只要Flink这边有checkpoint保存过状态就完全以Flink状态为准Kafka的group offset只是一个外部可见性的指标不是恢复的依据。只有第一次启动、没有任何历史状态时才会走启动模式StartupMode来决定起点EARLIEST、LATEST、GROUP_OFFSETS、SPECIFIC_OFFSETS、TIMESTAMP分别对应setStartFromEarliest、setStartFromLatest这些配置方法。这里有个容易翻车的场景如果作业曾经跑过但你在代码里把setStartFromLatest改成setStartFromEarliest它不会生效因为恢复逻辑是有状态用状态没状态才用启动模式。很多人以为改了启动模式就能从新位置消费重启后发现没变化其实就是这个原因。2.2 从KafkaConsumerThread到KafkaFetcher一个handover队列串起拉取与处理消费端的核心数据通路由两个线程组成。KafkaConsumerThread封装了原生org.apache.kafka.clients.consumer.KafkaConsumer它运行在独立线程里核心是一个poll循环不断调consumer.poll(timeout)拉取数据然后处理wakeup和可容忍异常。拉到的数据不能直接丢给下游算子因为它们不在同一个线程里。KafkaConsumerThread会把ConsumerRecords放进一个Handover队列这是一个带阻塞语义的交接器专门用于跨线程传递数据。KafkaFetcher运行在Flink的task线程里它执行runFetchLoop()阻塞在handover.pollNext()上拿到记录后逐条反序列化、过watermark逻辑再交给下游算子处理。这种双线程设计非常值得借鉴拉取的线程只管拉不碰任何业务逻辑处理的线程只管处理不碰Kafka客户端。两边用Handover解耦谁也不会因为对方的节奏拖累自己。更妙的是背压处理。当下游算子处理不动时KafkaFetcher自然停止从Handover取数据队列慢慢占满KafkaConsumerThread的下一步produce操作就会阻塞于是poll循环被卡住背压一路传回Kafka拉取端。整个链条不会像无界拉取那样把消息堆积在内存里也不会把数据直接怼给下游导致OOM。2.3 watermark的生成与多分区对齐的源码实现在消费端每个分区沿着数据流往上游注册自己的watermark。FlinkKafkaConsumerBase封装了AssignerWithPeriodicWatermarks或AssignerWithPunctuatedWatermarks的逻辑但真正的状态维护在KafkaPartitionState里每个分区有独立的timestamp和watermark字段。当多个分区被同一个并行子任务消费时最终这个并行子任务向上游发射的watermark是所有活跃分区里最小的那个。这是Flink取最小值保证不违反乱序的原则在消费端的直接体现。真正麻烦的是如果其中一个分区长时间没有新数据它的watermark会一直卡在旧值导致整个并行子任务的watermark被它拖住窗口不触发、side output不输出。源码里的解决办法是idle检测机制通过SimpleConsumerThread或KafkaFetcher中的idle时间逻辑把超过阈值没有数据的分区标记为idle从watermark对齐集合里剔除。我建议你在读这段源码时顺手把minWatermarkMark和idle的分区管理逻辑画个时序图这个理解了后面遇到某个分区的数据晚到导致窗口不触发的问题时你一眼就能判断是idle阈值没配还是对应分区真的没数据。3. 动态分区发现与checkpoint两条容易漏掉的源码细节3.1 分区发现线程做了什么分区发现是Kafka Connector源码里最容易被忽略、但生产环境最有用的机制。FlinkKafkaConsumerBase.run()方法启动时会创建KafkaPartitionDiscoverer它会利用Kafka客户端的管理接口拉取当前topic的完整分区列表。动态发现需要显式开启就是在传给FlinkKafkaConsumer的Properties里设置flink.partition-discovery.interval-millis这个参数比如Properties props new Properties(); props.setProperty(bootstrap.servers, localhost:9092); props.setProperty(group.id, flink-group); props.setProperty(flink.partition-discovery.interval-millis, 60000);如果不设置这个参数默认值是Long.MAX_VALUE也就是说分区发现只在作业启动时做一次后续Kafka新增的分区永远不会被消费。很多人把topic扩了分区后发现Flink作业没有任何反应其实就是没开这个参数数据全积压在新增分区里积压告警追到你脸上你才想起来。新发现的分区怎么决定起始offset分两种情况如果作业是从checkpoint恢复的新分区不在历史状态里没有已经保存的offset此时会走启动模式比如你配置的是setStartFromLatest那新分区就从最新的offset开始如果是作业首次启动同样按启动模式。这里要特别提醒新分区不会从当前消费进度继续走因为它本来就没有消费进度这个逻辑是刻在源码里的不是你用setStartFromLatest能控制的。3.2 checkpoint保存与offset提交的对应关系消费端的状态保存接口是snapshotState()在checkpoint触发时它会遍历当前所有KafkaTopicPartitionState把每个分区的当前offset下一条待消费位置写入状态后端。由于只是记录offset和分区编号这个状态非常小即使消费几千个分区单个checkpoint的状态通常也只有几百KB到一两MB所以你完全不用为这个状态大小焦虑。紧接着是notifyCheckpointComplete()如果配置了setCommitOffsetsOnCheckpoints(true)它会把已经保存的offset主动提交给Kafka的__consumer_offsets主题也就是Kafka的consumer group提交。这一步的作用是让Kafka侧看到Flink的消费进度让一些基于group offset的监控工具比如Kafka的消费组延迟监控能显示准确数值。但如果没开这个配置Flink的Exactly-Once语义并不会受损因为Flink恢复数据靠的是自己的状态不是Kafka的offset。读到这段源码时我最大的感悟是Flink把容错用的进度和暴露给外部的进度分得很清楚。前者是内部状态后者是可选的外部同步。很多人把这两件事混为一谈才会在排查Kafka消费组监控时被误导以为Flink没提交offset就是数据丢了其实根本不是一回事。4. 生产者链路从send到两阶段提交的Exactly-Once实现4.1 FlinkKafkaProducer的事务模型生产者链路的核心类是FlinkKafkaProducerBase它继承自Flink的TwoPhaseCommitSinkFunction这个抽象类把两阶段提交的骨架搭好了子类只需要实现beginTransaction、preCommit、commit、abort几个方法。FlinkKafkaProducerBase内部维护一个FlinkKafkaInternalProducer这是对原生Kafka生产者的包装增加了对事务状态的控制能力。每次checkpoint到来sink会调用beginTransaction()开启一个新事务数据写入过程中sink把这些数据缓存在当前事务的buffer里checkpoint完成前调用preCommit()把缓冲的数据flush到Kafka但还不提交事务等整个checkpoint都成功后再异步回调commit()提交事务。如果有任何一环失败就调用abort()回滚事务。Semantic枚举控制了这个流程的严格程度我整理了一个对照表Semantic底层行为适用场景注意点EXACTLY_ONCE开启Kafka事务两阶段提交端到端精确一次对数据一致性要求极高的金融、计数场景KAFKA服务的transaction.state.log配置需要正常且生产者事务超时时间要合理设置AT_LEAST_ONCE不加事务失败重放时可能重复写入大多数日志收集、监控指标场景下游要做幂等或去重否则会重复数据NONE不开启任何事务甚至不做等待对延迟极其敏感、允许丢失/重复的试验场景异常恢复时序无法保证生产环境不建议4.2 两阶段提交在Flink里如何被调度两阶段提交能在Flink里跑起来靠的是TwoPhaseCommitSinkFunction与checkpoint机制的配合。正常情况下checkpoint barrier流经sink算子时sink会先把已经收到的数据flush出去然后执行preCommit即把当前Kafka事务内的所有数据请求发送到broker并等待持久化但事务本身还是open状态。只有等JobManager确认整个checkpoint的各个环节都成功完成后才会调用notifyCheckpointComplete()在这个回调里再执行commitTransaction()让事务真正生效。这段逻辑看进去之后你会发现一个很实用的配置技巧Kafka事务是有超时时间的。如果你没有主动设置transaction.timeout.msKafka默认的事务超时可能是几十秒而你的checkpoint间隔如果设成了几分钟那preCommit的flush可能还没等到checkpoint completeKafka那边就自动把事务回滚了作业日志里会报TimeoutException。正确做法是把transaction.timeout.ms调大大于checkpoint间隔同时还得小等于Kafka broker端的transaction.max.timeout.ms否则服务端会直接拒绝。4.3 事务ID的生成与恢复容错事务ID是两阶段提交能实现容错的关键。FlinkKafkaProducerBase里事务ID的生成方式是transactionIdPrefix - subtaskIndex - checkpointId。为什么要带checkpointId因为每个checkpoint对应一个新事务checkpoint编号保证了同一子任务的不同周期事务ID不重复。作业失败重启并从checkpoint恢复时TwoPhaseCommitSinkFunction会走到恢复分支。这时如果上次的事务已经完成了preCommit但没commit新起来的事务会尝试从Kafka的__transaction_state主题里找回那个事务的状态决定是继续提交还是中止。这个设计的作用是避免两种异常情况一种是事务悬在已写入但未提交的状态另一种是事务被错误地重复提交。事务ID的幂等性至关重要。如果两个作业或者两次恢复用了相同的事务ID后开的事务会把前一个事务覆盖掉轻则丢数据重则整个事务协调器报冲突。这也就是为什么Flink要求每个作业的transactionIdPrefix要足够唯一尤其是多个作业共用同一个Kafka集群的时候别用默认前缀硬扛。4.4 三种Semantic在源码上的差异很多人以为Semantic只是一个简单的开关但源码里它们的路径完全不同。Semantic.NONE直接走简单的send路径不调用beginTransaction生产者失败恢复的语义完全不保证Semantic.AT_LEAST_ONCE也是直接send但会在checkpoint complete之后等待所有记录都被确认保证不丢但有重复只有Semantic.EXACTLY_ONCE真正走完整的两阶段提交路径beginTransaction→ 写入 →preCommit→ checkpoint complete →commit。实际使用中如果你把sink.setSemantic(Semantic.EXACTLY_ONCE)和setSemantic(Semantic.AT_LEAST_ONCE)在同样异常场景下对比你会看到完全不同的行为EXACTLY_ONCE模式下失败重放后Kafka broker侧不会出现重复记录AT_LEAST_ONCE模式下重放必然重复。这个差异在源码里就是有没有初始化事务的区别。5. 用源码知识解真实故障从连接器异常到火焰图热点5.1 连接器报错的排查链路连接器异常是群里的高频问题比如flink的jdbc连接器异常这类很多人一上来就贴日志问怎么办。我的习惯是先用读Kafka Connector源码时建立的方法论去解也就是先看报错类属于哪条链路再定位是配置问题、环境问题还是代码问题。拿Kafka Connector最常见的三个报错来说第一个是TimeoutException when trying to commit transaction。这个报错几乎都指向事务超时配置你去看FlinkKafkaProducerBase里的commitTransaction这段异常就是调用Kafka事务的commit时抛的。你优先检查transaction.timeout.ms和checkpoint.interval之间的关系再把Kafka broker端的transaction.max.timeout.ms查一下基本都能解决。第二个是OffsetOutOfRangeException。这个报错发生在fetcher拉取数据时说明Flink保存的下一条offset已经超出了Kafka分区当前保留的范围比如offset过期、日志被清理掉。源码里fetcher对这类异常会做周期性重试和状态更新所以作业通常不会直接挂掉但会出现某个分区一直消费不到新数据的状况。解决思路是确认是否有足够长的日志保留时间或者手动重置该分区的起点。第三个是UnknownTopicOrPartitionException。这个常见于topic被删了重建或者写代码时topic写错了。源码里分区发现器或fetch线程会把异常抛出来导致作业不断重启。排查时就去看KafkaConsumerThread.run()里对KafkaException的处理逻辑很多异常会被标记为可容忍并继续循环只有少部分会真正把任务搞崩。5.2 火焰图定位性能热点我之前在一次性能排查里用过Flink火焰图当时作业的CPU使用率很高但一直不知道热点在哪。后来我盯着KafkaConsumerThread.run和KafkaFetcher.runFetchLoop这两个方法看发现火焰图里KafkaConsumerThread.poll占的百分比异常高说明问题出在拉取端本身要么是Kafka客户端版本有坑要么是反序列化太重拖住了处理线程。如果火焰图里KafkaFetcher.runFetchLoop占大头那说明反序列化和下游emit才是瓶颈此时调拉取的批次大小反而没用你应该去看KafkaDeserializationSchema的实现是不是在每条记录上做了同步的、昂贵的外部调用。这类定位思路本质上是把源码的模块图变成你分析问题的地图。没有读过源码的人拿到火焰图只会觉得哪里都烫读过源码的人看到类名和方法名立刻能在脑子里找到对应的位置和配置项几分钟就能收敛问题范围。5.3 老Connector的局限与新KafkaSource的演进Flink从1.14开始推荐用KafkaSource老的FlinkKafkaConsumer虽然还能用但已经被标记为废弃。新老架构最核心的差异是老Connector把分区发现、恢复、拉取逻辑分散在FlinkKafkaConsumerBase、KafkaConsumerThread、KafkaFetcher等多个类里状态结构和线程模型都比较复杂新KafkaSource基于Source API重写每个Kafka分区作为一个split来管理KafkaPartitionSplitReader负责拉取KafkaRecordEmitter负责发记录状态管理集中在KafkaSourceEnumState里逻辑清楚很多。我的建议是如果你想深入学习先读老Connector因为它把各种机制暴露得最直观理解成本其实更低如果做新项目直接用KafkaSource省心。读完旧代码再对照新代码你会发现新版很多设计就是针对旧版的痛点来改的这种版本设计对比带给你的理解比单独读任何一版都深刻。最后分享一个我自己的习惯读完Kafka Connector的源码别急着合上电脑去把社区里关于这段设计的讨论翻出来看看那里记录了设计者踩过的坑和权衡过程。你会发现很多看起来别扭的代码都是历史原因和现实约束共同作用的结果比如老Connector里为什么有那么多针对不同Kafka版本的子类就是因为上游客户端API变化太快Flink只能跟着打补丁。把这段历史补上你的源码阅读才算是真的闭环。以上是我在实际阅读和排查中的一些体会希望能给你一点参考。
返回列表