Flink 源码与运行原理

流式计算 / 源码阅读 / LESSON 02

流数据端到端链路——Source → Transform → Sink 的全源码穿越

沿 Source、Transform、Sink 的端到端链路,理解流记录如何在 Operator Chain 与网络层中移动。

阅读时间
120 分钟
学习路径
Flink 源码与运行原理
内容来源
Flink 深度笔记

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


下一步