预计阅读时间: 120 分钟 前置阅读: doc-00(架构概念), doc-01(如何调度) 下一次阅读: doc-03(Checkpoint), doc-05(Window)
1. StreamTask 执行循环—Flink 的心跳
1.1 架构定位
StreamTask 是所有流处理 Task 的基类。核心执行循环 processInput() 在一个 while 循环中不断读取输入、处理、输出结果。
1.2 Mailbox 模型详解
Flink 的 StreamTask 使用 Mailbox 模型 实现单线程协作式多任务:
// MailboxProcessor.java — Mailbox 执行循环
public void runMailboxLoop() {
// 外层循环: task 的生命周期循环
while (isNextLoopPossible()) {
// 内层1: 处理所有待处理的邮件 (非阻塞)
processMail(localMailbox, false);
// 邮件包括:
// - Checkpoint Trigger (高优先级)
// - Timer 回调
// - Checkpoint Completion 通知
// - Operator Event
if (isNextLoopPossible()) {
// 内层2: 执行默认动作 (processInput)
mailboxDefaultAction.runDefaultAction(mailboxController);
}
}
}
// 邮件处理的优先级:
// MAX_PRIORITY: Checkpoint 通知
// MIN_PRIORITY: 其他所有邮件
为什么需要 Mailbox 模型?
- 传统方案:在
processElement()中间被中断(如 Checkpoint Trigger)→ 需要处理中断、保存中间状态 - Mailbox 方案:每个 StreamElement 处理完毕后检查邮件队列 → 无需处理中断 → 代码简洁且安全
1.3 StreamTask 完整调用链
// StreamTask.java — 生命周期
// ===== 1. restore() 阶段 =====
public final void restore() throws Exception {
// 创建 OperatorChain (算子链)
operatorChain = new RegularOperatorChain<>(this, recordWriter);
// 初始化 (task-specific, 如创建 StreamInputProcessor)
init();
// 恢复状态和 InputGate
restoreStateAndGates();
// → StateBackend.restore() → 从 Checkpoint 恢复
// → InputGate 恢复 (异步)
}
// ===== 2. invoke() 阶段 =====
public final void invoke() throws Exception {
// 启动 BufferDebloater
scheduleBufferDebloater();
// 进入 Mailbox 循环 → 反复调用 processInput()
mailboxProcessor.runMailboxLoop();
// 循环退出 → afterInvoke()
}
// ===== 3. processInput() — 核心数据处理 =====
protected void processInput(MailboxDefaultAction.Controller controller) {
while (true) {
// 从 InputGate 读一个 StreamElement
DataInputStatus status = inputProcessor.processInput();
switch (status) {
case MORE_AVAILABLE:
// 有更多数据可用 → 立即返回 (再次进入 processInput)
return;
case NOTHING_AVAILABLE:
// 无数据 → 检查反压/空闲
if (!recordWriter.isAvailable()) {
// 输出端无缓冲 → 反压!
controller.suspendDefaultAction(recordWriter.getAvailableFuture());
return;
}
// 空闲 → 等待新数据
controller.suspendDefaultAction(inputProcessor.getAvailableFuture());
return;
case END_OF_INPUT:
// 所有输入已完成 → 通知 EndOfData
endData(DRAIN);
controller.allActionsCompleted();
return;
case STOPPED:
// Task 停止
controller.allActionsCompleted();
return;
}
}
}
1.4 processInput 循环中的事件分发
// StreamInputProcessor.java — 事件分发
protected void processInput() {
// 从 InputGate 读取下一个 BufferOrEvent
BufferOrEvent bufferOrEvent = inputGate.getNext().orElse(null);
if (bufferOrEvent.isBuffer()) {
// 普通数据 → 反序列化 → StreamRecord
StreamRecord record = deserializeRecord(bufferOrEvent.getBuffer());
operator.processElement(record); // 交给算子链的第一个算子
} else if (bufferOrEvent.getEvent() instanceof Watermark) {
// Watermark → 触发 Timer 或 Window
operator.processWatermark((Watermark) bufferOrEvent.getEvent());
} else if (bufferOrEvent.getEvent() instanceof CheckpointBarrier) {
// Barrier → 触发 Checkpoint 快照
operator.processBarrier((CheckpointBarrier) bufferOrEvent.getEvent());
} else if (bufferOrEvent.getEvent() instanceof LatencyMarker) {
// 延迟追踪
operator.processLatencyMarker((LatencyMarker) bufferOrEvent.getEvent());
}
}
2. Source 算子—数据从哪里来
2.1 FLIP-27 Source 架构
Source (用户 API)
│
├── SplitEnumerator (JM 端运行, 单例)
│ ├── start() // 初始化
│ │ └── 发现 Splits (Kafka: listPartitions)
│ ├── handleSplitRequest() // 处理 Split 请求
│ │ └── 将 Partition 分配给请求的 Reader
│ ├── addSplitsBack() // Task 失败 → Split 退回
│ └── snapshotState() // Checkpoint: 保存 Split 分配
│
└── SourceReader (TM 端运行, 每个并行度 1 个)
├── pollNext() // 被 StreamTask 循环调用
│ └── 从 KafkaConsumer 读取 → 反序列化 → StreamRecord
├── addSplits() // 接收分配的新 Split
│ └── 创建/更新 KafkaConsumer 的 Partition 分配
└── snapshotState() // Checkpoint: 保存 Kafka offset
└── offset → StateHandle → AcknowledgeCheckpoint
2.2 Kafka Source 具体流程
KafkaSourceReader
→ KafkaConsumer.poll(timeout) // poll() 返回 ConsumerRecords
→ 遍历 ConsumerRecords:
→ deserialize(record) // KafkaDeserializationSchema
→ output.collect(streamRecord) // 进入算子链
→ 如果有记录: return MORE_AVAILABLE
→ 如果无记录: return NOTHING_AVAILABLE
→ StreamTask 根据返回值决定: 立即循环(有数据) vs 等待(无数据)
Offset 管理 (Checkpoint 时):
KafkaSourceReader.snapshotState()
→ 构建 offset 映射: Partition → currentOffset
→ 写入 OperatorState → StateHandle → Ack
→ Checkpoint 完成后 → Kafka Committer 提交 offset 到 Kafka
2.3 FLIP-27 vs 旧 SourceFunction
| 维度 | 旧 SourceFunction | FLIP-27 Source API |
|---|---|---|
| Split 发现 | 与读取耦合 | 分离 (SplitEnumerator 独立) |
| Checkpoint | 手动保存状态 | 框架自动管理 |
| 动态 Partition | 需要手动重启 | 自动发现和分配 |
| 流批统一 | 不支持 | 同一 Source 实现流批 |
| Offset 管理 | SourceContext.getCheckpointLock() | 框架自动管理 |
| 线程模型 | SourceFunction 独占线程 | StreamTask 主线程 |
3. Operator Chain—算子链优化
3.1 链化的收益
不链化 (每个算子独立 Task):
Source Task → [序列化] → [Netty] → [反序列化] → Map Task
每个 Record → 4 次边界跨越:
函数调用 + 序列化 + TCP 传输 + 反序列化 + 函数调用
链化 (Source + Map 在一个 Task):
Source.processElement() → Map.processElement()
同线程内函数调用传递 → 零序列化、零网络
每个 Record → 1 次函数调用
性能对比:
- 不链化:延迟增加 ~50-200μs (序列化 + TCP)
- 链化:延迟增加 ~0.1μs (函数调用)
3.2 链化条件(源码)
// StreamingJobGraphGenerator.java — 链化判定
private static boolean isChainable(StreamNode upStream, StreamNode downStream) {
// 1. 并行度相同
if (upStream.getParallelism() != downStream.getParallelism()) return false;
// 2. SlotSharingGroup 相同
if (!upStream.getSlotSharingGroup().equals(downStream.getSlotSharingGroup()))
return false;
// 3. 数据分发为 FORWARD
for (StreamEdge edge : upStream.getOutEdgesInOrder()) {
if (edge.getTargetId() == downStream.getId()) {
if (edge.getPartitioner() != ForwardPartitioner.INSTANCE) return false;
}
}
// 4. ChainingStrategy 允许
if (upStream.getChainingStrategy() == ChainingStrategy.NEVER) return false;
if (downStream.getChainingStrategy() == ChainingStrategy.NEVER) return false;
// 5. 算子是同一类型 (不能把 Source 和 Operator 链化? No, 可以)
// Source → Map → Filter → Sink 可以全部链化
return true;
}
3.3 链化后的执行
// OperatorChain.java — 链式调用
public void processElement(StreamRecord record) throws Exception {
// 链头: Source.processElement() 产生 StreamRecord
// → headOperator.processElement(record)
// 链中: 自动传播
// Map.processElement() 中调用 output.collect(record)
// → output 自动路由到链中的下一个算子
// → Filter.processElement(record)
// 链尾: Sink.processElement(record) 写入外部
// 在同一个调用栈中完成 (深度优先链式调用)
}
4. Shuffle—RecordWriter → ResultPartition → InputGate
4.1 数据分发模式
| 模式 | 分区器 | 目标选择 | 触发 |
|---|---|---|---|
| FORWARD | ForwardPartitioner | 固定 1:1 | 默认 (同并行度) |
| HASH | KeyGroupStreamPartitioner | hash(key) % parallelism | keyBy() |
| REBALANCE | RebalancePartitioner | Round-robin | rebalance() |
| BROADCAST | BroadcastPartitioner | 所有下游 | broadcast() |
| RESCALE | RescalePartitioner | Round-robin (仅本地) | rescale() |
| GLOBAL | GlobalPartitioner | 全部发给 Subtask 0 | global() |
| CUSTOM | CustomPartitioner | 用户自定义 | partitionCustom() |
4.2 RecordWriter.emit() — 数据写入
// RecordWriter.java — emit 流程
public void emit(T record, int targetSubpartition) throws IOException {
// 1. 序列化
serializer.setPosition(0);
serializeRecord(serializer, record);
// 格式: [Length:4B][Data:Variable]
// 2. 选择目标子分区 (由分区器决定)
int target = partitioner.partition(record, numberOfSubpartitions);
// 3. 写入子分区
targetPartitionWriter.emit(serializer.getSharedBuffer(), target);
// 4. 如果需要立即 flush
if (flushAlways) {
targetPartitionWriter.flush(target);
}
}
// ===== 事件广播 (Watermark, Barrier) =====
public void broadcastEvent(AbstractEvent event) {
for (int i = 0; i < numberOfSubpartitions; i++) {
// 将事件序列化并写入所有子分区
// 确保所有下游 Subtask 都收到相同的 Watermark / Barrier
targetPartitionWriter.emitEvent(event, i);
}
}
5. Sink 算子—数据到哪里去
5.1 FLIP-143 Sink API 架构
Sink (用户 API)
│
├── SinkWriter (TM 端, 每个并行度 1 个)
│ ├── write(element, context) // 写数据 (如追加到批次)
│ ├── prepareCommit() // 预提交 (如 flush 批次)
│ │ └── 返回 Committable 列表
│ └── snapshotState() // Checkpoint 时调用
│
└── Committer (JM 端, 全局 1 个)
├── commit(committables) // 真正提交 (如 Kafka 事务 commit)
│ └── Checkpoint 完成后调用
└── retryRecovery(committables) // 恢复时重试
5.2 两阶段提交(Kafka Sink 示例)
Phase 1 (Pre-commit, 在 Checkpoint 快照阶段):
1. SinkWriter.prepareCommit()
→ KafkaProducer.flush() // 刷出所有未发送数据
→ beginTransaction() 的事务信息存入 Committable
2. SinkWriter.snapshotState() ← Checkpoint 快照
→ 保存 Committable 到 OperatorState
3. CheckpointCoordinator 收集所有 StateHandle
→ 所有算子的快照完成 → Checkpoint Complete
Phase 2 (Commit, 在 Checkpoint 完成后):
4. NotifyCheckpointComplete(checkpointId)
→ Committer.commit(committables)
→ KafkaProducer.commitTransaction() // 真正提交!
→ Offset 被标记为 Committed
Phase 3 (Abort, 在 Checkpoint 失败时):
5. SinkWriter.abort(committables)
→ KafkaProducer.abortTransaction()
为什么两阶段提交保证 Exactly-Once?
- Checkpoint 中同时包含了 Source Offset 和 Sink 事务信息
- 恢复时:Source 回到 Checkpoint 中的 Offset + Sink 重放/重试未完成的事务
- 要么事务提交(Checkpoint 成功),要么事务回滚(Checkpoint 失败)→ 没有中间状态
6. StreamElement 类型体系
// StreamElement — 流元素的继承体系
abstract class StreamElement {
boolean isRecord(); // 是否为数据记录
boolean isWatermark(); // 是否为 Watermark
boolean isStreamStatus(); // 是否为流状态
boolean isLatencyMarker(); // 是否为延迟标记
}
StreamRecord<T> extends StreamElement {
T value; // 用户数据
long timestamp; // 事件时间戳
boolean hasTimestamp; // 是否有有效时间戳
}
Watermark extends StreamElement {
long timestamp; // 水印值
// "时间 ≤ timestamp 的数据均已到达"
}
StreamStatus extends StreamElement {
StreamStatus IDLE; // Source 空闲 (Watermark 不推进)
StreamStatus ACTIVE; // Source 活跃
}
LatencyMarker extends StreamElement {
long timestamp; // 创建时间
// 追踪端到端延迟
}
7. 源码导航(完整版)
| 文件 | 关键类/方法 | 职责 |
|---|---|---|
flink-streaming-java/.../tasks/StreamTask.java | processInput(), invoke(), restore(), afterInvoke() | 流处理 Task 基类 |
flink-streaming-java/.../tasks/mailbox/MailboxProcessor.java | runMailboxLoop(), processMail() | Mailbox 执行模型 |
flink-streaming-java/.../tasks/StreamInputProcessor.java | processInput() | 事件分发 |
flink-streaming-java/.../streamrecord/StreamRecord.java | — | 流元素: value + timestamp |
flink-streaming-java/.../streamrecord/StreamElement.java | — | 流事件基类 |
flink-core/.../connector/Source.java | createReader(), createEnumerator() (FLIP-27) | Source API |
flink-core/.../connector/Sink.java | createWriter(), createCommitter() (FLIP-143) | Sink API |
flink-streaming-java/.../operators/ChainingStrategy.java | — | 链化策略枚举 |
flink-optimizer/.../StreamingJobGraphGenerator.java | createJobGraph(), isChainable() | JobGraph 生成 + 链化 |
flink-runtime/.../io/network/api/writer/RecordWriter.java | emit(), broadcastEmit(), randomEmit() | 记录写入器 |
flink-runtime/.../taskmanager/Task.java | doRun() | Task 包装 |
flink-streaming-java/.../io/RecordWriterOutput.java | collect(), emitWatermark() | 流事件输出 |
8. 常见问题 / 面试题
Q1: StreamTask 的 Mailbox 模型和 Actor 模型的异同?
A: 相似点:(1) 单线程处理;(2) 通过消息(邮件)传递控制信息。不同点:(1) Mailbox 是协作式(yield control at safe points),Actor 是抢占式(any time);(2) Flink 的 Mailbox 专为流处理优化——在每处理一条 Record 后才检查邮件,确保原子性;(3) Mailbox 的默认动作(processInput)可以暂时挂起并等待异步 future。
Q2: Operator Chain 什么时候不生效?关键是哪个条件最容易触发断链?
A: 最常见的断链原因是 keyBy() — 它改变了数据分发模式(FORWARD → HASH)。KeyBy 引入了跨网络的 Shuffle(需要 hash(key) 决定目标 Subtask),不同 Subtask 可能在不同 TM 上 → 必须断链。其他原因:显式 disableChaining()、不同 SlotSharingGroup、改变并行度。
Q3: Barrier 和数据记录的顺序如何保证?
A: Barrier 和数据记录在同一个 Channel(ResultSubpartition → InputChannel)中顺序传输。TCP 保证字节流的顺序性,而 Barrier 在 PipelinedSubpartition 中被放入优先级队列——但优先级队列只是 Barrier 优先于其他数据被读取,不会打乱 Barrier 与数据在同一 Channel 内的相对顺序。Barrier Buffer 在收到时阻塞上游(isBlocked=true),确保 Barrier 后面的数据不会"越过"Barrier 被处理。
Q4: FLIP-27 Source 的 SplitEnumerator 为什么在 JM 端运行?
A: SplitEnumerator 需要全局视图来决定 Partition → Subtask 的映射。如果每个 Subtask 自行决定,两个 Subtask 可能会消费同一个 Partition(重复消费)。JM 端的单例 SplitEnumerator 保证分配的全局唯一性。
Q5: 两阶段提交和普通 Sink 的本质区别?
A: 普通 Sink(如 KinesisAtLeastOnceProducer)在 Checkpoint 前就直接写外部了。Job 失败后 Checkpoint 回滚 → Source 回退 Offset → 重新消费 → Sink 再次写入 → 重复数据(At-Least-Once)。两阶段提交在 Checkpoint 成功后才 commit() → 失败时事务自动回滚 → 不产生重复数据(Exactly-Once)。
下一步
- Checkpoint 详解: doc-03-checkpoint-and-savepoint.md
- Window & Timer: doc-05-window-and-timer.md
- 网络栈: doc-09-network-stack.md