flink kafka-connector0.10版本分析,与1.4版本中kafka11对比
引言
官方文档在出了1.4之后特意发表了一篇blog,通过以下这两个条件实现了真正意义上的exactly once语义
- kafka producer0.11的事务性
- two phase commit protocol
我们先看0.10版本的kafka-connector的行为逻辑.
Kafka-Connector10
kafkaConsumer10
FlinkKafkaConsumerBase
这个抽象类实现了CheckpointedFunction, 这个接口的描述:
1 | * This is the core interface for <i>stateful transformation functions</i>, meaning functions |
这个接口中主要要去做两件事情:
1 | //每一次做checkpoint的时候被调用 |
初始化函数的调用时机是在open之前的:
1 | //StreamTask.java |
在初始化的函数中提供了一个FunctionSnapshotContext1
void snapshotState(FunctionSnapshotContext context) throws Exception;
让你既可以注册一个KeyedStateStore,也可以注册一个OperatorStateStore
1 | public interface ManagedInitializationContext { |
我们可以看到kafka10是怎么利用这个CheckpointedFunction来管理记录内部offset的呢?
initializeState
1 | //初始化过程 |
那么我们看到恢复或者初始化的时候将
1 | //将从状态中获取到的列表赋予给消费的列表 |
snapshot
1 |
|
offsetCommitMode
如上述代码中的offsetCommitMode,主要有以下几种
1 | DISABLED, |
这个配置主要改变了commitOffset回kafka的时机. 首先在snapshot的时候会将对应的checkpointId和相应的offset的列表放入pendingOffsetsToCommit, 在checkpoint完成后回调notifyCheckpointComplete,这里面主要完成了offset的commit工作。
1 | if (offsetCommitMode == OffsetCommitMode.ON_CHECKPOINTS) { |
消费partition分配问题
在从上次点恢复的情况下是直接从state中获取应该读取哪一个partition,offset。如果并发度改变了会做出什么样的反馈呢?会正确做出rescale吗
第一次进行读取的时候会初始化处当前task(并发度)所需要订阅的partition
1 | protected static Map<KafkaTopicPartition, Long> initializeSubscribedPartitionsToStartOffsets( |
1 | public static int assign(KafkaTopicPartition partition, int numParallelSubtasks) { |
消费模型kafka10
在完成partition订阅之后,就要开始真正的run方法了,FlinkKafkaConsumer也是实现自SouceFunction,因此主要的逻辑也都是在run方法中实现。 主要逻辑:
kafkaConsumerThread和Kafka10Fetcher通过Handover交互,我觉得这段代码写的很不错,可以好好学习下。可以形象的比作在接力跑:kafkaConsumerThread通过真正的消费线程消费放入一个HandOver,再由kafkaFetcher去poll,完成整个消费过程。
1 | // we need only do work, if we actually have partitions assigned |
1 | // Handover的描述, Handover代码可以再好好学习下。 |
KafkaFetcher
1 | public void runFetchLoop() throws Exception { |
1 | protected void emitRecordWithTimestamp( |
在设置了kafkaTimestampassigner之后就会进行一个定时任务向下游发送watermark,值为所有partition维护的最小值:
1 | public void onProcessingTime(long timestamp) throws Exception { |
KafkaConsumerThread
1 | if (records == null) { |
主要做的工作就是从consumer消费数据塞入handover,等待拉取
Handover
桥接模式
1 | public class KafkaConsumerCallBridge { |
解决08 09 10 版本的api不兼容问题
Kafkaproducer10
###initializeState
什么都不做
snapshot
1 |
|
这里涉及到一个flushOnCheckPoint的问题,再调用producer.flush期间,producer会将所有没写入的,在buffer中的数据刷盘,然后调用commitCallBack,这就保证了ckpt之后数据不会丢的问题。
主要工作方法:
1 | public void invoke(IN next) throws Exception { |
问题
kafka恢复状态直接从状态中去获取了之前保存的partition和offset,但是如果是扩容partition的场景就不会从新的Partition消费 issue:FLINK-8869
flink内部维护了offset,为什么向kafka提交的时候还需要在checkpoint之后再提交而不是定时提交就算了?
因为虽然从checkpoint点恢复的时候不需要从kafka broker获取消费点的位置了,但是如果是应用重启消费上次消费到的点的数据,这个offset就是flink向kafka提交的,放在checkpoint完成后去做的好处就是让应用即使不是从上个点恢复的,也能够从kafka消费正确的offset点。
如果在新的checkpoint没打之前任务失败了,重新从上次的offset点消费的话下游数据是不是重复了?
是的,因为有一部分数据经过处理已经sink出去了,因此才需要0.11的一致性语义
以上代码:
kafka-connector0.10 来源于release1.3.2
kafka-connector0.11 来源于release1.4.0