Flink internals

STREAM PROCESSING / SOURCE READING / LESSON 10

TaskManager memory model

Break down heap, direct, managed, and network memory in Flink's TaskManager model.

Reading
90 min
Track
Flink internals
Source
Chinese source notes

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

预计阅读时间: 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 的聚合结果)。


下一步