Flink 源码与运行原理

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

Flink 内存模型深度解析

拆解 TaskManager 内存模型,区分堆内、直接内存、托管内存与网络内存的用途和约束。

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

预计阅读时间: 90 分钟 前置阅读: doc-00(架构) 下一次阅读: doc-11(RM 调度)


1. 两套内存模型:JobManager vs TaskManager

1.1 JobManager 内存模型

JM 内存结构相对简单,主要存储元数据(不处理数据流):

┌──────────────────────────────────────────────────────┐
│              JobManager JVM 进程                      │
├──────────────────────────────────────────────────────┤
│  JVM Heap                                            │
│  │  Flink Framework Objects                          │
│  │  ├── Flink UI (WebServer)                         │
│  │  │   └── 默认约 50MB (可配置)                      │
│  │  ├── ExecutionGraph (JobGraph→物理图)              │
│  │  │   └── 每个 ExecutionVertex ~1-5KB              │
│  │  │   └── 1000 并行度 → ~10MB                       │
│  │  ├── CheckpointCoordinator                        │
│  │  │   └── PendingCheckpoint 元数据                  │
│  │  │   └── CompletedCheckpointStore (历史 Checkpoint) │
│  │  │   └── 1000 个算子 × 3 保留 Checkpoint → ~50MB   │
│  │  ├── JobGraph / StreamGraph 缓存                   │
│  │  ├── RPC (Akka/Netty) 缓冲区                      │
│  │  └── BlobServer (JAR 分发)                         │
│  │                                                    │
│  │  JM Heap 配置:                                     │
│  │  jobmanager.memory.heap.size (默认: 无, 从 process.size 推算) │
│  │                                                     │
│  Off-Heap Memory                                      │
│  │  ├── Metaspace (类元数据, 默认 256MB)               │
│  │  │   └── jobmanager.memory.metaspace.size           │
│  │  ├── Direct Memory (Netty I/O)                     │
│  │  │   └── 默认 ~128MB                                │
│  │  └── JVM Overhead (线程栈, GC, Native)              │
│  │      └── jobmanager.memory.jvm-overhead.fraction (0.1) │
│  │                                                      │
│  Total Process Memory:                                 │
│  jobmanager.memory.process.size (默认 1600MB)           │
└──────────────────────────────────────────────────────┘

JM 配置参数

参数默认说明
jobmanager.memory.process.size1600MBJM 进程总内存
jobmanager.memory.heap.sizeJM Heap 大小
jobmanager.memory.off-heap.size128MBJM Off-Heap 大小
jobmanager.memory.metaspace.size256MBJM Metaspace 大小
jobmanager.memory.jvm-overhead.fraction0.1JM JVM Overhead 比例

JM 内存估算公式

JM 内存 ≈ 1.6GB 默认 + Job 规模增量

Job 规模增量:
  - 小 Job (<100 并行度): 默认 1.6GB 足够
  - 中 Job (100-500 并行度): 2-4GB
  - 大 Job (500-2000 并行度): 4-8GB
  - 巨 Job (2000+ 并行度): 8-16GB

1.2 TaskManager 内存模型(完整版)

┌──────────────────────────────────────────────────────────┐
│                TaskManager JVM 进程                       │
│                                                          │
│  ┌────────────────────────────────────────────────────┐ │
│  │              Framework Heap                        │ │
│  │  - Flink Framework 自身对象                         │ │
│  │  - 默认 ~128MB                                      │ │
│  │  - taskmanager.memory.framework.heap.size           │ │
│  ├────────────────────────────────────────────────────┤ │
│  │              Task Heap (用户代码)                    │ │
│  │  - 用户 Function / Operator 对象                    │ │
│  │  - HashMapStateBackend 的状态数据                   │ │
│  │  - taskmanager.memory.task.heap.size                │ │
│  └────────────────────────────────────────────────────┘ │
│                                                          │
│  ┌────────────────────────────────────────────────────┐ │
│  │            Managed Memory (Off-Heap)                │ │
│  │  ┌──────────────────────────────────────────────┐  │ │
│  │  │ RocksDB: MemTable (write buffer)              │  │ │
│  │  │ RocksDB: Block Cache (read cache)             │  │ │
│  │  │ Sort-Merge Join: 排序缓冲区                    │  │ │
│  │  │ Hash Join: Build-side hash table              │  │ │
│  │  │ Python UDF: Python 进程内存                    │  │ │
│  │  │ taskmanager.memory.managed.size                │  │ │
│  │  └──────────────────────────────────────────────┘  │ │
│  ├────────────────────────────────────────────────────┤ │
│  │            Network Memory (Off-Heap)               │ │
│  │  - NetworkBuffer (每个 32KB)                       │ │
│  │  - 默认: flink.size × 0.1 (10%)                   │ │
│  │  - taskmanager.memory.network.fraction             │ │
│  │  - 最少 64MB                                       │ │
│  ├────────────────────────────────────────────────────┤ │
│  │            Direct Memory (框架 I/O)                │ │
│  │  - Netty 的 DirectBuffer                          │ │
│  │  - JNI 调用分配                                    │ │
│  │  - = taskmanager.memory.framework.off-heap.size   │ │
│  └────────────────────────────────────────────────────┘ │
│                                                          │
│  ┌────────────────────────────────────────────────────┐ │
│  │              Metaspace (类元数据)                   │ │
│  │  - 默认 256MB                                       │ │
│  │  - taskmanager.memory.metaspace.size                │ │
│  └────────────────────────────────────────────────────┘ │
│                                                          │
│  ┌────────────────────────────────────────────────────┐ │
│  │            JVM Overhead (不可控 Native 内存)         │ │
│  │  - 线程栈 (默认 1MB/线程)                           │ │
│  │  - GC 内部结构                                      │ │
│  │  - Code Cache                                       │ │
│  │  - taskmanager.memory.jvm-overhead.fraction         │ │
│  └────────────────────────────────────────────────────┘ │
└──────────────────────────────────────────────────────────┘

1.3 内存计算器公式

Flink 内存计算公式(由 TaskExecutorMemoryConfiguration 计算):

给定 taskmanager.memory.process.size = 4096MB (TM 总内存):

1. Metaspace = taskmanager.memory.metaspace.size = 256MB
2. JVM Overhead = process.size * jvm-overhead.fraction = 4096 * 0.1 = 410MB
3. Flink Total = process.size - Metaspace - JVM Overhead = 4096 - 256 - 410 = 3430MB

   Flink Total 内部:
   a. Framework Heap = 128MB (默认)
   b. Task Heap = ? (推导)
   c. Framework Off-Heap = 128MB (默认)
   d. Network Memory = Flink Total * network.fraction = 3430 * 0.1 = 343MB
   e. Managed Memory = Flink Total * managed.fraction = 3430 * 0.4 = 1372MB
   f. Task Heap = Flink Total - Framework Heap - Task Heap...

   或者,如果显式指定 managed.size:
     Task Heap + Network = Flink Total - Framework - Managed

2. MemoryManager — Managed Memory 管理

2.1 MemoryManager 类结构

// MemoryManager.java — Managed Memory 的池化管理
public class MemoryManager {
    private final long memorySize;           // 总 Managed Memory 大小
    private final int pageSize;              // 内存页大小 (默认 32KB)
    private final int numberOfPages;         // 总页数
    private final UnsafeMemoryBudget budget; // 原子预算

    // 空闲段池 (ArrayDeque, 无锁+锁保护)
    private final ArrayDeque<MemorySegment> freeSegments;

    // 已分配段: owner → Set<MemorySegment>
    private final ConcurrentHashMap<Object, Set<MemorySegment>> allocatedSegments;

    // 保留内存 (不分配段,只扣除预算): owner → size
    private final ConcurrentHashMap<Object, Long> reservedMemory;
}

2.2 两种分配模式

// ===== 模式 1: 基于页的分配 (用于算子内存) =====
public List<MemorySegment> allocatePages(Object owner, int numPages)
    throws MemoryAllocationException {

    // 1. 从 UnsafeMemoryBudget 预留
    if (!budget.reserve(numPages * pageSize)) {
        throw new MemoryAllocationException("Not enough managed memory");
    }

    // 2. 分配 MemorySegment (堆外 Unsafe 内存)
    List<MemorySegment> segments = new ArrayList<>(numPages);
    for (int i = 0; i < numPages; i++) {
        segments.add(allocateOffHeapUnsafeMemory(pageSize));
    }

    // 3. 记账
    allocatedSegments.computeIfAbsent(owner, k -> new HashSet<>())
        .addAll(segments);

    return segments;
}

// ===== 模式 2: 保留式分配 (用于 RocksDB) =====
// RocksDB 自己管理内存,Flink 只做预算控制
public boolean reserveMemory(Object owner, long size) {
    if (!budget.reserve(size)) return false;
    reservedMemory.merge(owner, size, Long::sum);
    return true;
}

2.3 SharedMemoryResources — RocksDB 的 Managed Memory 集成

// MemoryManager.java — RocksDB 获取 Managed Memory
public <T extends AutoCloseable> OpaqueMemoryResource<T>
    getSharedMemoryResourceForManagedMemory(
        String resourceType,       // "rocksdb"
        SharedMemoryResource.Initializer<T> initializer,
        double fraction) {

    // 1. 计算分配大小
    long numBytes = (long) (memorySize * fraction);

    // 2. 预留预算
    if (!budget.reserve(numBytes)) {
        throw new MemoryAllocationException(...);
    }

    // 3. 调用初始化器 (创建 RocksDB WriteBufferManager 或 Block Cache)
    T resource = initializer.initialize(numBytes);

    // 4. 返回带引用计数和释放器的 OpaqueMemoryResource
    return new OpaqueMemoryResource<>(resource, numBytes, () -> {
        resource.close();
        budget.release(numBytes);
    });
}

// 使用示例 (RocksDB Backend 内部):
// 为 RocksDB 的 WriteBufferManager 分配 Managed Memory
OpaqueMemoryResource<WriteBufferManager> writeBufferResource =
    memoryManager.getSharedMemoryResourceForManagedMemory(
        "rocksdb-write-buffer",
        size -> new WriteBufferManager(size, cache),
        writeBufferFraction  // 0.5 of managed memory
    );

// 为 RocksDB Block Cache 分配
OpaqueMemoryResource<Cache> blockCacheResource =
    memoryManager.getSharedMemoryResourceForManagedMemory(
        "rocksdb-block-cache",
        size -> new LRUCache(size),
        blockCacheFraction   // 0.3 of managed memory
    );

3. Network Memory — BufferPool 层次结构

3.1 三层缓冲池架构

Global: NetworkBufferPool (单例,TM 级别)
  ┌────────────────────────────────────────────────┐
  │ totalNumberOfMemorySegments = N (每个 32KB)      │
  │ availableMemorySegments: Queue<MemorySegment>   │
  │ requestMemorySegment() → 返回一个 Segment        │
  │ recycle(MemorySegment) → 回收                    │
  │                                                  │
  │ 总网络内存 = taskmanager.memory.network.fraction  │
  │            × Flink Total Memory                  │
  └────────────────────────────────────────────────┘
                │ 分配各 N 个给每个 Pool
    ┌───────────┼───────────┬───────────────┐
    ▼           ▼           ▼               ▼
 LocalBufferPool    LocalBufferPool    LocalBufferPool ...
 (ResultPartition)  (InputGate)        (InputGate)
  min: 2             min: 2             min: 2
  max: 无限制         max: 无限制         max: 无限制

每个 LocalBufferPool:
  - currentPoolSize: 当前持有的 Segment 数量
  - availableMemorySegments: 空闲的 Segment
  - numRequired: 最少需要的 Segment 数
  - maxOverdraftBuffersPerGate: 透支缓冲上限

3.2 NetworkBuffer 的分配和回收

// LocalBufferPool.java — 申请 NetworkBuffer
public MemorySegment requestMemorySegment(int targetChannel) {
    // 1. 本地池获取
    MemorySegment segment = availableMemorySegments.poll();
    if (segment != null) return segment;

    // 2. 从全局池获取(增加 currentPoolSize)
    segment = networkBufferPool.requestMemorySegment();
    if (segment != null) {
        currentPoolSize++;
        return segment;
    }

    // 3. 透支缓冲(超出 currentPoolSize 但有限制)
    if (canUseOverdraftBuffer(targetChannel)) {
        segment = networkBufferPool.requestMemorySegment();
        if (segment != null) return segment;
    }

    // 4. 真的没内存了 → 反压!
    availabilityHelper.resetUnavailable();
    return null;
}

// NetworkBufferPool.java — 全局池分配
public MemorySegment requestMemorySegment() {
    synchronized (availableMemorySegments) {
        return availableMemorySegments.poll();  // 可能为 null
    }
}

// NetworkBuffer 被 InputStream 消费后回收:
// 1. NetworkBuffer.recycleBuffer()
// 2. → BufferRecycler.recycle(MemorySegment)
// 3. → LocalBufferPool.recycle(MemorySegment)
// 4. → availableMemorySegments.offer(segment)
// 5. → availabilityHelper.available()  // 唤醒等待者

4. OOM 场景与排查

4.1 常见 OOM 类型

OOM 类型原因排查修复
OutOfMemoryError: Java heap spaceTask Heap 不足查看 Heap 使用趋势增大 taskmanager.memory.task.heap.size
OutOfMemoryError: Direct buffer memoryNetwork Memory + Direct Memory 超限查看 Direct Memory 使用增大 network.fraction, 减小 managed.fraction
OutOfMemoryError: Metaspace类加载过多jstat -gc <pid>增大 taskmanager.memory.metaspace.size
OutOfMemoryError: unable to create new native thread线程数超限ulimit -u / 线程数减小并行度或增大 JVM Overhead
OutOfMemoryError: GC overhead limit exceededGC 频繁但回收少GC 日志增大 Heap 或切换到 RocksDB Backend
RocksDB write stallManaged Memory 不足RocksDB 指标增大 taskmanager.memory.managed.size

4.2 OOM 排查流程

1. 确认 OOM 类型:
   - 查看 TM/JM 日志: grep "OutOfMemoryError"
   - 确认是 Heap / Direct / Metaspace / Native

2. 查看内存配置:
   - Web UI → TaskManager → Metrics → Memory
   - heapUsed / heapCommitted / heapMax
   - directUsed / directTotal

3. 确认内存泄漏 or 合理不足:
   - 内存持续增长 → 泄漏: 排查用户代码中未清理的 State
   - 内存稳定高位 → 不足: 增大对应内存区

4. 火焰图确认:
   - async-profiler 采样 Heap 分配路径
   - 找出占用内存最多的对象

5. 修复:
   - 调整内存配置
   - 优化用户代码(减少状态、设置 TTL)
   - 切换 StateBackend(Heap → RocksDB)

5. 多 Topic 接入的内存影响分析

5.1 影响维度总览

场景: 1 个 Flink Job 消费 3 个 Kafka Topic (TopicA, TopicB, TopicC)

              TopicA ──┐
              TopicB ──┼──→ KafkaSource(并行度=8) → Window → Sink
              TopicC ──┘

每个 Topic 的 Partition 被均匀分配到 8 个 Source Subtask

5.2 内存影响

内存区单 Topic多 Topic (3个)增量
Task HeapSource Reader 对象 + 反序列化缓冲3x Source Reader + 3x 缓冲~2-3x
Network Memory总吞吐决定,与 Topic 数无关(相同总吞吐)相同总吞吐 → 相同 Network Memory0%
Managed MemoryRocksDB State: 3 个 Topic 的 offset 存储3x State 数量~3x State
Direct MemoryNetty I/O多 Channel → 更多 DirectBuffer轻微增加
Checkpoint 内存3 个 Topic 的 Offset 快照增加了 State Handle 数量5-10%

5.3 关键风险

风险 1: State 膨胀
  TopicA: 低流量 → 状态稳定
  TopicB: 高流量 + 大 Key 空间 → 状态持续增长
  TopicC: 突发流量 → 状态波动
  → 所有 Topic 的 State 在同一个 Window 算子中 → 整体 Managed Memory 压力叠加

风险 2: Checkpoint 被慢 Topic 拖累
  TopicA 数据多 → Barrier 对齐慢
  所有 Source 的 Barrier 必须全部到达才能完成 Checkpoint
  → 慢 Topic 决定 Checkpoint 速度

风险 3: 反压传播
  TopicB 流量过大 → 下游算子处理不过来
  → 反压传导回所有 Source
  → TopicA 和 TopicC 也被迫停止消费

风险 4: Watermark 停滞
  TopicC 长时间无数据 → 对应 Partition 不推进 Watermark
  → 多输入算子的 Watermark = min(all inputs)
  → 窗口不触发 → 状态堆积

5.4 推荐配置

# 多 Topic 场景推荐配置
taskmanager.memory.managed.size: 适当增大(考虑所有 Topic 的总状态)
taskmanager.memory.network.fraction: 0.15(多 Channel 场景增加网络缓冲)
taskmanager.numberOfTaskSlots: 控制 Slot 数(避免过多 Source 并发)

# Checkpoint 配置
execution.checkpointing.interval: 适当增大(避免 Checkpoint 过于频繁)
execution.checkpointing.unaligned: true(慢 Topic 对齐问题)
execution.checkpointing.aligned-checkpoint-timeout: 30s(自适应切换)

# Watermark 配置
WatermarkStrategy.forBoundedOutOfOrderness(...)
    .withIdleness(Duration.ofMinutes(1))  // 空闲 Topic 不阻塞 Watermark

6. 关键配置速查表

配置默认说明调优建议
taskmanager.memory.process.size1728MBTM 进程总内存生产环境 4-32GB
taskmanager.memory.flink.sizeFlink 框架总内存约 = process.size - Metaspace - Overhead
taskmanager.memory.task.heap.sizeTask Heap用户代码内存,无 States 时默认 ~512MB
taskmanager.memory.managed.sizeManaged MemoryRocksDB 场景建议占 flink.size 的 40-50%
taskmanager.memory.managed.fraction0.4Managed Memory 占比RocksDB → 0.4~0.5, HashMap → 0.0
taskmanager.memory.network.fraction0.1Network Memory 占比高吞吐场景可增大到 0.15-0.2
taskmanager.memory.network.min64MBNetwork Memory 最小值
taskmanager.memory.network.max1GBNetwork Memory 最大值
taskmanager.memory.jvm-overhead.fraction0.1JVM Overhead 占比线程多时增大到 0.15-0.2
taskmanager.memory.jvm-overhead.min192MBJVM Overhead 最小值
taskmanager.memory.jvm-overhead.max1GBJVM Overhead 最大值
taskmanager.memory.framework.heap.size128MBFramework Heap
taskmanager.memory.framework.off-heap.size128MBFramework Off-Heap
taskmanager.memory.metaspace.size256MBMetaspace大量 UDF 时可增大

7. 源码导航(完整版)

文件关键类/方法职责
flink-runtime/.../memory/MemoryManager.javaallocatePages(), reserveMemory(), getSharedMemoryResourceForManagedMemory()Managed Memory 管理
flink-runtime/.../memory/MemorySegment.javaallocateOffHeapUnsafeMemory(), free()内存段
flink-runtime/.../memory/TaskExecutorMemoryConfiguration.javaTM 内存配置计算
flink-runtime/.../memory/UnsafeMemoryBudget.javareserve(), release(), getAvailableMemorySize()内存预算
flink-runtime/.../network/buffer/NetworkBuffer.java网络缓冲
flink-runtime/.../network/buffer/NetworkBufferPool.javacreateBufferPool(), requestMemorySegment()全局网络缓冲池
flink-runtime/.../network/buffer/LocalBufferPool.javarequestMemorySegment(), recycle()本地缓冲池
flink-runtime/.../network/buffer/BufferPool.javagetAvailableFuture()反压信号
flink-state-backends/flink-statebackend-rocksdb/.../RocksDBOptionsFactory.javacreateDBOptions(), createColumnFamilyOptions()RocksDB 内存配置

8. 常见问题 / 面试题

Q1: Managed Memory 和 Direct Memory 的核心区别?

A: Managed Memory 由 Flink 的 MemoryManager 管理,用于算子级计算(RocksDB 写缓存、排序缓冲区、Hash Join)。Direct Memory 由 JVM 的 java.nio.DirectByteBuffer 管理,通过 -XX:MaxDirectMemorySize 控制上限,用于 Netty 的网络 I/O 和 JNI 调用。两者都在 Off-Heap,但管理方式和使用者不同。

Q2: 为什么 Network Memory 默认只占 10%?

A: 10% 在大多数场景足够。NetworkBuffer 只用于缓冲传输中的数据(flight 中的网络包),一旦消费者消费完成就回收。实际需要的 Network Memory 与 "并发传输缓冲数 × 32KB" 相关,与总内存大小不是线性关系。但高吞吐、多 Channel 场景可增大到 15%-20%。

Q3: 如何判断是增大 Managed Memory 还是 Task Heap?

A:

  • HashMapStateBackend → 状态在 Task Heap 中 → 增大 task.heap.size
  • RocksDB StateBackend → 状态在 Managed Memory (Off-Heap) → 增大 managed.size
  • 反例:把 RocksDB 的 managed.fraction 减小并增大 task.heap.size(错误——Task Heap 空着但 Managed Memory 不够)

Q4: JM 挂了会导致内存不足吗?什么情况需要增大 JM 内存?

A: 通常 JM 1.6GB 足够。需要增大的场景:

  • 大并行度 Job(2000+ 并行度),ExecutionGraph 元数据占用大
  • Checkpoint 保留数量多(> 10),Checkpoint 元数据占用大
  • 大量 Operator(> 100),JobGraph 元数据占用大
  • HaServices 使用 ZooKeeper 且有大量 Checkpoint 元数据

Q5: 多 Topic 接入时,如果每个 Topic 的 Schema 不同(不同反序列化器),对内存有什么额外影响?

A: 主要是 Task Heap 增加——每个 Topic 需要独立的 Source Reader 实例、独立的 KafkaDeserializationSchema 实例和独立的序列化缓冲。此外,每个 Topic 的数据可能进入不同的下游算子分支 → 更多的 Operator Chain → 更多的 Task → 更多的 Framework Heap 开销。如果所有 Topic 的数据最终汇聚到同一个 Window → 只有 Managed Memory 的增量(RocksDB 中存储更多 Key 的聚合结果)。


下一步