Kafka JDBC Source Connector最后读取数据库行的存储及实现问询
我想从实现层面了解Kafka Connect JDBC Source Connector是怎么追踪已经读取并发布的最后一条数据库行的。
我查过Confluent官方文档,里面只提到功能但没说具体实现:
Kafka Connect tracks the latest record it retrieved from each table, so it can start in the correct location on the next iteration (or in case of a crash). The source connector uses this functionality to only get updated rows from a table (or from the output of a custom query) on each iteration. Several modes are supported, each of which differs in how modified rows are detected.
我也翻了连接器的源码,但还没找到具体实现细节。我自己猜测它可能把最后发布的行信息存在数据库/键值存储里,或者在发布的Kafka主题里加元数据,方便下一轮轮询时获取?
核心实现逻辑:依赖Kafka Connect的偏移量存储框架
JDBC Source Connector本身并不独立实现追踪存储,而是完全依托Kafka Connect内置的偏移量管理机制来记录已读取的最后一行位置,具体细节如下:
偏移量的存储位置:
分布式模式下,偏移量默认存在Kafka内部的__consumer_offsets主题;本地独立模式下,偏移量会存在本地文件(路径由offset.storage.file.filename配置)。Connect框架负责偏移量的持久化和恢复,连接器只需要提交追踪值即可。不同同步模式下的偏移量内容:
根据配置的同步模式(mode参数),连接器提交的偏移量数据会不同:incrementing模式:存储最后读取的自增主键字段值,下一轮轮询只查询主键大于该值的行timestamp模式:存储最后读取的时间戳字段值,下一轮轮询只查询时间戳大于该值的行timestamp+incrementing模式:同时存储自增主键和时间戳值,避免时间戳重复导致的漏读/重复读snapshot模式:全量同步完成后,会自动切换到对应增量模式并记录初始偏移量
源码中的关键实现:
在JdbcSourceTask类的poll()方法中,每次轮询完成后会调用context.offsetStorageWriter().offsets(...)提交偏移量;而偏移量的具体值是从TableQuerier的实现类(比如处理增量+时间戳模式的TimestampIncrementingTableQuerier)中获取的当前轮询到的最大追踪字段值。
纠正你的猜测
它既不会把追踪信息存在业务数据库,也不会在发布的Kafka主题中添加元数据——所有追踪逻辑都通过Kafka Connect的标准偏移量框架完成,这也是Kafka生态组件的统一设计模式。
内容的提问来源于stack exchange,提问作者Tony

