Flink internals

STREAM PROCESSING / SOURCE READING / LESSON 02

End-to-end stream dataflow

Follow records from source to transform and sink through operators, task chains, and the network layer.

Reading
120 min
Track
Flink internals
Source
Chinese source notes

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

预计阅读时间: 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

维度旧 SourceFunctionFLIP-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 数据分发模式

模式分区器目标选择触发
FORWARDForwardPartitioner固定 1:1默认 (同并行度)
HASHKeyGroupStreamPartitionerhash(key) % parallelismkeyBy()
REBALANCERebalancePartitionerRound-robinrebalance()
BROADCASTBroadcastPartitioner所有下游broadcast()
RESCALERescalePartitionerRound-robin (仅本地)rescale()
GLOBALGlobalPartitioner全部发给 Subtask 0global()
CUSTOMCustomPartitioner用户自定义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.javaprocessInput(), invoke(), restore(), afterInvoke()流处理 Task 基类
flink-streaming-java/.../tasks/mailbox/MailboxProcessor.javarunMailboxLoop(), processMail()Mailbox 执行模型
flink-streaming-java/.../tasks/StreamInputProcessor.javaprocessInput()事件分发
flink-streaming-java/.../streamrecord/StreamRecord.java流元素: value + timestamp
flink-streaming-java/.../streamrecord/StreamElement.java流事件基类
flink-core/.../connector/Source.javacreateReader(), createEnumerator() (FLIP-27)Source API
flink-core/.../connector/Sink.javacreateWriter(), createCommitter() (FLIP-143)Sink API
flink-streaming-java/.../operators/ChainingStrategy.java链化策略枚举
flink-optimizer/.../StreamingJobGraphGenerator.javacreateJobGraph(), isChainable()JobGraph 生成 + 链化
flink-runtime/.../io/network/api/writer/RecordWriter.javaemit(), broadcastEmit(), randomEmit()记录写入器
flink-runtime/.../taskmanager/Task.javadoRun()Task 包装
flink-streaming-java/.../io/RecordWriterOutput.javacollect(), 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)。


下一步