The source notes for this track are currently maintained in Chinese.

1. FLIP-27 Source API 深度
1.1 架构回顾
Source<T, SplitT, EnumChkT>
│
├── SplitEnumerator<SplitT, EnumChkT> (JM 端, 单例)
│ ├── start() 初始化发现 Splits
│ ├── handleSplitRequest(host, id) 响应 Reader 的 Split 请求
│ ├── addSplitsBack(splits, id) Task 失败 → 退回 Split
│ ├── addReader(readerId) 注册新 Reader
│ └── snapshotState(checkpointId) Checkpoint: 保存 Split 分配
│
└── SourceReader<T, SplitT> (TM 端, 每个 Subtask 1 个)
├── pollNext(output) StreamTask 循环调用
├── addSplits(splits) 接收新分配 Split
├── handleSourceEvents(event) 处理枚举器发来的事件
└── snapshotState(checkpointId) Checkpoint: 保存 offset
1.2 Kafka Source 数据读取流程
// KafkaSourceReader.java — 核心读取
public InputStatus pollNext(ReaderOutput<T> output) {
// 1. 从 KafkaConsumer 拉取数据
ConsumerRecords<byte[], byte[]> records =
kafkaConsumer.poll(Duration.ofMillis(pollTimeout));
// 2. 逐条反序列化
for (ConsumerRecord<byte[], byte[]> record : records) {
T value = deserializationSchema.deserialize(record);
// 3. 输出 StreamRecord
output.collect(value, record.timestamp());
}
// 4. 有数据 → 立即再次循环
// 无数据 → 等待 (StreamTask 的 NOTHING_AVAILABLE)
return records.isEmpty() ? InputStatus.NOTHING_AVAILABLE : InputStatus.MORE_AVAILABLE;
}
1.3 Kafka Source 的 Offset 管理
Offset 管理三步走 (Source 端):
1. SourceReader.snapshotState() → 记录当前 offset
2. CheckpointCoordinator 收集所有 Ack → Checkpoint Complete
3. Kafka Committer: NotifyCheckpointComplete → KafkaConsumer.commitSync(offsets)
为什么不在 SourceReader 中直接 commit offset?
→ 如果 Reader poll() 后 + 数据处理前 Commit → 数据丢失!
→ 必须等 CheckpointComplete 后才 commit → Exactly-Once
2. FLIP-143 Sink API 与两阶段提交
2.1 Sink API 架构
Sink<T>
│
├── SinkWriter<T, CommT> (TM 端)
│ ├── write(element, context) 写数据
│ ├── prepareCommit() 预提交 (flush + 生成 Committable)
│ │ └── 返回 List<Committable>
│ └── snapshotState(checkpointId) Checkpoint 时调用
│
├── Committer<CommT> (JM 端, 全局)
│ ├── commit(committables) 真正提交 (Checkpoint 完成后)
│ └── retryRecovery(committables) 恢复时重试
│
└── GlobalCommitter<CommT, GlobalCommT> (JM 端, 可选)
└── combine(committables) 全局合并 (如合并小文件)
2.2 Kafka Sink 两阶段提交完整流程
Phase 1: Pre-commit (在 Checkpoint 快照时)
1. KafkaWriter.write(record)
→ KafkaProducer.send(record) // 发送但仍未提交事务
2. KafkaWriter.prepareCommit()
→ KafkaProducer.flush() // 等所有发送完成
→ return new KafkaCommittable(txnId) // 返回事务 ID
3. KafkaWriter.snapshotState()
→ 保存 KafkaCommittable 到 OperatorState
→ CheckpointCoordinator 收集所有 StateHandle
Phase 2: Checkpoint (Barrier 在 DAG 中传播)
4. CheckpointCoordinator.completePendingCheckpoint()
→ 所有 Subtask 的 Ack 都已到达
→ Checkpoint #N Complete
Phase 3: Commit (在 Checkpoint 完成通知中)
5. NotifyCheckpointComplete(N)
→ KafkaCommitter.commit(txnId) // 真正提交!
→ KafkaProducer.commitTransaction(txnId)
→ offset 被标记为 Committed
Phase 4: Abort (在 Checkpoint 失败时)
6. abortCheckpointOnCheckpointFailure(N)
→ KafkaWriter.abort(txnId)
→ KafkaProducer.abortTransaction(txnId)
2.3 两阶段提交 vs 普通 Sink
| 维度 | 普通 Sink | 两阶段提交 |
|---|---|---|
| 写外部时机 | 实时 (每条 Record) | 实时 (但事务未提交) |
| Checkpoint 失败 | 已写外部数据无法回滚 → 重复 | 事务自动回滚 → 无重复 |
| 一致性 | At-Least-Once | Exactly-Once |
| 代价 | 低 | 事务开销 + Checkpoint 对齐 |
| Support | 所有 Sink | 支持事务的外部系统 (Kafka, JDBC 等) |
3. 自定义 Connector 开发流程
1. 实现 Source<T, SplitT, EnumChkT>
→ getSplitSerializer() / createReader() / createEnumerator()
2. 实现 SplitEnumerator
→ start() → 发现 Split (如文件列表、Partition 列表)
→ handleSplitRequest() → 分配 Split (对等分配 or Hash 分配)
→ snapshotState() → 保存 Split 分配状态
3. 实现 SourceReader
→ pollNext() → 从分配的 Split 读取数据 → output.collect()
→ addSplits() → 接收新 Split (初始化读取器)
→ snapshotState() → 保存读取位置 (offset)
4. (可选) 实现 Split 状态序列化器
→ SimpleVersionedSerializer<SplitT>
4. 源码导航(完整版)
| 文件 | 关键类/方法 | 职责 |
|---|---|---|
flink-core/.../connector/Source.java | createReader(), createEnumerator() | Source 接口 |
flink-core/.../connector/Sink.java | createWriter(), createCommitter() | Sink 接口 |
flink-connectors/.../kafka/source/KafkaSource.java | — | Kafka Source |
flink-connectors/.../kafka/source/KafkaSourceReader.java | pollNext() | Kafka Reader |
flink-connectors/.../kafka/source/KafkaSourceEnumerator.java | start(), handleSplitRequest() | Kafka 枚举器 |
flink-connectors/.../kafka/sink/KafkaSink.java | — | Kafka Sink (两阶段提交) |
flink-connectors/.../kafka/sink/KafkaWriter.java | write(), prepareCommit(), snapshotState() | Kafka Writer |
5. 常见问题 / 面试题
Q1: FLIP-27 解决了旧 SourceFunction 的什么问题?
A: (1) 分离 Split 发现(Enumerator)和读取(Reader)→ 动态 Partition 发现;(2) 框架自动管理 Offset → 用户不需要手动 getCheckpointLock();(3) 流批统一——同一 Source 实现支持批 (有界流) 和流 (无界流);(4) 线程模型简化——Reader 在 StreamTask 主线程中运行,不需要单独的 Source Thread。
Q2: 两阶段提交如何保证 Exactly-Once?Checkpoint 失败了怎么办?
A: Sink 端在 Checkpoint 前只做 prepareCommit(flush 但不 commit),Checkpoint 完成后才 commit。Checkpoint 失败 → Checkpoint 回滚 → Source 回退到旧 Offset → 重新消费 → Sink 的旧事务自动回滚 → 新事务覆盖 → 无重复。
Q3: 自定义 Connector 时必须实现什么?最少的工作量?
A: 最少工作量:(1) Source<MyType, MySplit, MyEnumeratorState> 的 3 个工厂方法;(2) SplitEnumerator<MySplit, MyEnumeratorState> 的 split 发现 + 分配逻辑;(3) SourceReader<MyType, MySplit> 的 pollNext() 读取逻辑 + snapshotState() 位置保存。如果 Split 结构简单(如 1 个 Partition = 1 个 Split),约 200-300 行代码。