
Flink 之所以能在流计算领域站住,靠的不只是引擎本身,更靠连接器(Connector)生态这块基石。作业的「源头」和「去处」都长在连接器上:Kafka 取数、MySQL CDC 抓变更、Hudi/Iceberg 落湖、ClickHouse 写仓……连接器质量直接决定了一条流能不能「读得准、写得对、看得见」。
最近两个进入主干的改动,恰好把「读得准」和「看得见」两块补上了:一个是关联 FLINK-33484 的 PR,修复了 Flink Kafka 连接器中 Offset 提交(offset commit) 相关的问题,提升了数据消费的准确性与一致性;另一个是关联 FLINK-40397 的 PR,为 MySQL Binlog 读取器新增了 currentBinlogPositionLag 指标,用来监控当前 Binlog 的消费延迟。
一个管「提交到 Kafka 的 offset 到底准不准」,一个管「MySQL 变更追到哪了、离最新差多少」。二者一内一外,共同指向连接器成熟度的同一道题:消费位点的正确性与可观测性。下面分两大部分,把底层机制和源码设计讲透。
先说清楚 Kafka 连接器的位点语义。Flink 1.12 之后的新 KafkaSource(基于 FLIP-27 Source API)彻底重构了旧 FlinkKafkaConsumer 的模型。核心角色是:
SourceReader;KafkaSourceReader):在 TaskManager 侧,每个 reader 持有若干 partition 的 KafkaPartitionSplit,用内部的 KafkaPartitionSplitReader 去 poll 数据;__consumer_offsets,是一致性的关键。关键认知:新 KafkaSource 默认关闭 Kafka 自身的
enable.auto.commit。位点提交完全由 Flink 接管,且只在 Checkpoint 完成时触发。这意味着:一个位点被「提交到 Kafka」的充要条件是「它所在的 Checkpoint 已经成功完成」。这是 exactly-once / at-least-once 语义的底座。
KafkaSourceReader 内部用一张「待提交位点表」串起 snapshot 与 checkpoint-complete 两阶段:
public class KafkaSourceReader<T>
extends SourceReaderBase<T, KafkaPartitionSplit, KafkaPartitionSplitState, KafkaSourceEnumState> {
/** checkpointId -> 该检查点对应的待提交位点(每个 partition 一个 offset)*/
private final Map<Long, Map<TopicPartition, OffsetAndMetadata>> pendingOffsetsToCommit;
@Override
public List<KafkaPartitionSplitState> snapshotState(long checkpointId) {
// 1. 先让父类把各 split 的进度快照(供故障恢复用)
List<KafkaPartitionSplitState> states = super.snapshotState(checkpointId);
// 2. 记录「当前已消费到的位点」作为该 checkpoint 的待提交位点
pendingOffsetsToCommit.put(checkpointId, currentConsumedOffsets());
return states;
}
@Override
public void notifyCheckpointComplete(long checkpointId) {
// 3. Checkpoint 成功后,把对应位点提交给 Kafka
Map<TopicPartition, OffsetAndMetadata> offsets = pendingOffsetsToCommit.remove(checkpointId);
if (offsets != null) {
consumer.commitSync(offsets); // 写回 __consumer_offsets
}
}
}注意 currentConsumedOffsets() 取的是「已经消费到」的位点。这一步的取值方式,正是 FLINK-33484 修复的关键所在。

老实现里,currentConsumedOffsets() 依赖 split reader 内部「最近一次 poll 下来的 offset」。问题就出在「poll 下来」和「真正被下游处理」之间存在时间差与并发窗口:
SplitReader 内部被异步更新,snapshot 时刻取到的值,可能比实际已发送下游的位点更靠前(把还没 emit 的数据也算进去了),也可能因 split 被回收/重新分配而拿到过期的位点;FLINK-33484 的修复思路是:把「待提交位点」锚定到「该 split 最后一个真正 emit 给下游的记录的位点」,而不是「poll 下来但未处理」的位点;同时对 split 回收、rescale 重分配等边界做了位点合并与去重,保证 pendingOffsetsToCommit 里每个 partition 的 offset 既不会超前、也不会漏算。

对照上图,修复前是「按 poll 位置提交 → 可能超前/落后」的长链路;修复后是「按 emit 位置提交 → 与实际处理严格对齐」的短链路。在 exactly-once 下游(如写数据库、发下游 Kafka)场景里,这一修直接决定了重启后数据是「恰好一次」还是「悄悄丢了半条」。
下面这张时序图把 JobManager(CheckpointCoordinator)、KafkaSourceReader、SplitReader、Kafka Broker 之间的关系画清楚。重点是两条:snapshot 时记录「已 emit 位点」,以及 notifyCheckpointComplete 时才真正 commitSync——提交动作永远晚于位点产生,这保证了提交的一定是已确认进检查点的进度。

源码层面的几个要点:
pendingOffsetsToCommit 以 checkpointId 为键,天然支持乱序完成的检查点(旧的提交不会覆盖新的);commitSync 而非 commitAsync,保证提交失败能被 Checkpoint 机制感知,不会「提交丢了但检查点说成功」;这里必须点清一个常被误解的点:新 KafkaSource 基于 Checkpoint 的位点提交,提供的是 at-least-once 语义(在 DeliveryGuarantee.NONE 下甚至不提交),而非端到端的 exactly-once。道理很直接——Checkpoint 成功意味着「位点已提交 + 下游状态已快照」;但「位点提交」和「下游真正落库」之间存在时间窗:
TwoPhaseCommitSinkFunction / 事务型 Sink)时,才能把「读—算—写」对齐成端到端 exactly-once。所以 FLINK-33484 的价值不在于「变出 exactly-once」,而在于先消灭掉更隐蔽的那一类错误:在只有 at-least-once 的前提下,提交位点超前会造成「数据凭空消失」,这比重复消费更难发现、更致命。修复把「重复」的上界钉死在「一次 Checkpoint 间隔内」,把「丢失」彻底归零。
作业并行度调整(rescale)时,partition split 会在 subtask 之间重新分配。旧路径下,一个 split 从 reader A 移交给 reader B 时,若 A 的 pendingOffsetsToCommit 里还残留该 split 的过期位点,就可能出现「B 用新位点、A 却把旧位点提交了」的竞争。修复后的 handover 协议要求:
SplitState,随 split 信封一并移交;一句话:位点永远跟着 split 走,而不是跟着线程/reader 实例走。这是把一致性从「大概率对」收敛到「结构上对」的关键设计。
讲完 Kafka 的「读得准」,再讲 MySQL CDC 的「看得见」。
Flink CDC 的 MySQL 连接器(flink-cdc)通过 Debezium 引擎读取 MySQL 的 binlog。Binlog 是 MySQL 的变更流,连接器像「从库」一样追着主库的 binlog 文件一路读。但 binlog 是无限流,你永远在「追」最新的位置——追得快不快,就是延迟指标要回答的问题。
MySQL 连接器的位点用 BinlogOffset 表示,核心字段是 fileName(如 mysql-bin.000012)+ position(文件内字节偏移)+ 可选 gtid:
public class BinlogOffset implements Comparable<BinlogOffset> {
private final String filename; // mysql-bin.000012
private final long position; // 该文件内的偏移
private final String gtidSet; // 可选 GTID
/** 比较两个 binlog 位点的先后,供延迟计算使用 */
public int compareTo(BinlogOffset o) { ... }
}读取器在消费 binlog 事件时,每处理完一个事件,就把「当前事件对应的 binlog 位点」记到 BinlogSplitState 里——这就是消费进度。
延迟 = 「MySQL 服务端当前最新 binlog 位点」−「连接器已消费到的位点」。前者靠周期性 SHOW MASTER STATUS 拿到,后者就是上面那个消费进度。
FLINK-40397 把这个差值做成了 Flink 指标系统里的一个 Gauge,名字叫 currentBinlogPositionLag:
public class MySqlBinlogSplitReader implements SplitReader<SourceRecord, BinlogSplitState> {
private BinlogOffset currentBinlogOffset; // 已消费到的位点
private BinlogOffset serverBinlogOffset; // 服务端最新位点(定期刷新)
@Override
public void handleEvent(SourceRecord record) {
// 每消费一个 binlog 事件,刷新当前位点
this.currentBinlogOffset = toBinlogOffset(record);
}
/** FLINK-40397 新增:将延迟注册为 Gauge */
private void registerMetrics(SourceReaderContext context) {
context.getMetricGroup().gauge(
"currentBinlogPositionLag",
() -> computeLag(currentBinlogOffset, serverBinlogOffset));
}
private long computeLag(BinlogOffset cur, BinlogOffset server) {
if (cur == null || server == null) return 0L;
// 用 binlog 文件序号 + 文件内偏移综合估算「差了多少」
return binlogDistance(server, cur);
}
}
整条链路是:MySQL 主库写 binlog → Debezium/MySqlSplitReader 追读 → 反序列化事件 emit 给下游 → 同时把当前位点喂给 currentBinlogPositionLag 这个 Gauge;而 Gauge 的另一端,是定时 SHOW MASTER STATUS 拿到的服务端最新位点。运维在 Grafana 上看到这条曲线的瞬间,就知道「这条 CDC 管道是不是在追、追得及不及时」。
没有它之前,MySQL CDC 延迟只能从「下游写入延迟」「Flink 反压」间接推断,定位不到「是源端读得慢,还是下游写得慢」。有了 currentBinlogPositionLag:

它的实现遵循 Flink 指标的惯用法:在 SourceReader 初始化时通过 SourceReaderContext.getMetricGroup() 拿到指标组,把一个 Supplier<Long> 注册成 Gauge;运行时 Flink 指标系统按采样周期去拉这个值,再上报给 Prometheus / JMX。读取器本身不感知上报细节,只负责算出「差多少」。
computeLag 看似只做减法,实则要处理「跨 binlog 文件」的比较。binlog 文件名形如 mysql-bin.000012,序号是单调递增的。一个朴素的「文件内偏移相减」在跨文件时会得到负数,必须先用文件名序号对齐:
private long binlogDistance(BinlogOffset server, BinlogOffset cur) {
// 1. 不同文件:用文件序号差 + 当前文件剩余量估算
int f1 = parseFileIndex(server.filename); // 12
int f2 = parseFileIndex(cur.filename); // 11
if (f1 != f2) {
// 粗略按「(f1-f2) 个文件 + 各自偏移」折算,单位字节
return (long)(f1 - f2) * AVG_BINLOG_FILE_BYTES
+ server.position + (AVG_BINLOG_FILE_BYTES - cur.position);
}
// 2. 同一文件:直接做偏移差
return server.position - cur.position;
}如果开启了 GTID 模式,则可退化为「已执行 GTID 集合 vs 服务端 GTID 集合」的事务数差,精度更高。FLINK-40397 落地的就是这个计算口径,让 currentBinlogPositionLag 在「刚切换 binlog 文件」这种边界时刻也不会出现跳变或负值。
拿到指标只是第一步。生产里通常这样用:
currentBinlogPositionLag 按 source 实例分组画曲线,配一条阈值线(如 50MB 或 30s 等效)做告警;关键认知:这个指标度量的是「源端追赶速度」,不是「端到端延迟」。它回答「连接器追 binlog 追得及不及时」,不回答「一条变更从发生到落湖用了多久」。后者还要叠加下游处理延迟,二者别混为一谈。
为了把两个改动放在一起看,先给一张对照表:
维度 | FLINK-33484(Kafka Offset 提交) | FLINK-40397(MySQL Binlog 延迟) |
|---|---|---|
连接器 | Kafka Source(FLIP-27) | Flink CDC MySQL Source |
解决痛点 | 提交的 offset 超前/落后,重启丢数或重复 | 源端消费延迟不可见,瓶颈难定位 |
核心改动 | 待提交位点锚定到「已 emit 位点」 | 新增 |
机制归类 | 消费位点的正确性 | 消费进度的可观测性 |
触达位置 |
|
|
故障影响 | 数据一致性(静默丢数) | 排障效率(黑盒变曲线) |
二者一个管「提交准不准」、一个管「进度看得见」,共同把连接器从「能跑通 demo」推进到「生产可放心依赖」。
回到开头那句话——连接器质量决定流能不能「读得准、写得对、看得见」。FLINK-33484 把 Kafka 的 offset 提交从「按 poll 位置(可能超前/落后)」钉成了「按 emit 位置(与实际处理严格对齐)」,消除了重启后静默丢数或重复消费的隐患;FLINK-40397 给 MySQL CDC 装上了 currentBinlogPositionLag 这双眼睛,让「追 binlog 追得及不及时」从黑盒变成可量化曲线。
二者一个管「提交准不准」、一个管「进度看得见」,共同把连接器从「能跑通 demo」推进到「生产可放心依赖」。
如果你也在用 Flink 接 Kafka 或做 MySQL CDC,建议顺着两条线自查:Kafka 侧,确认 enable.auto.commit=false 且消费位点确实由 Checkpoint 驱动,别让「提交位点」和「实际处理」之间留口子;CDC 侧,把 currentBinlogPositionLag 接进看板,把「源端追不上」和「下游写不动」两类问题一眼分开。这两处往往就是「同样一个数据不一致,别人定位在源、你定位三天还在下游日志里翻」的差距所在。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。