Flink internals

STREAM PROCESSING / SOURCE READING / LESSON 12

Connector framework and Source/Sink APIs

Study connector lifecycles, concurrency, and consistency through FLIP-27 Source and Sink APIs.

Reading
60 min
Track
Flink internals
Source
Chinese source notes

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

预计阅读时间: 60 分钟 前置阅读: doc-02(流数据链路) 下一次阅读: doc-13(恢复)


Connector 框架与 Source/Sink API 深度解析 图 01

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-OnceExactly-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.javacreateReader(), createEnumerator()Source 接口
flink-core/.../connector/Sink.javacreateWriter(), createCommitter()Sink 接口
flink-connectors/.../kafka/source/KafkaSource.javaKafka Source
flink-connectors/.../kafka/source/KafkaSourceReader.javapollNext()Kafka Reader
flink-connectors/.../kafka/source/KafkaSourceEnumerator.javastart(), handleSplitRequest()Kafka 枚举器
flink-connectors/.../kafka/sink/KafkaSink.javaKafka Sink (两阶段提交)
flink-connectors/.../kafka/sink/KafkaWriter.javawrite(), 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 行代码。


下一步