The source notes for this track are currently maintained in Chinese.
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.size | 1600MB | JM 进程总内存 |
jobmanager.memory.heap.size | — | JM Heap 大小 |
jobmanager.memory.off-heap.size | 128MB | JM Off-Heap 大小 |
jobmanager.memory.metaspace.size | 256MB | JM Metaspace 大小 |
jobmanager.memory.jvm-overhead.fraction | 0.1 | JM 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 space | Task Heap 不足 | 查看 Heap 使用趋势 | 增大 taskmanager.memory.task.heap.size |
OutOfMemoryError: Direct buffer memory | Network 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 exceeded | GC 频繁但回收少 | GC 日志 | 增大 Heap 或切换到 RocksDB Backend |
| RocksDB write stall | Managed 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 Heap | Source Reader 对象 + 反序列化缓冲 | 3x Source Reader + 3x 缓冲 | ~2-3x |
| Network Memory | 总吞吐决定,与 Topic 数无关(相同总吞吐) | 相同总吞吐 → 相同 Network Memory | 0% |
| Managed Memory | RocksDB State: 3 个 Topic 的 offset 存储 | 3x State 数量 | ~3x State |
| Direct Memory | Netty 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.size | 1728MB | TM 进程总内存 | 生产环境 4-32GB |
taskmanager.memory.flink.size | — | Flink 框架总内存 | 约 = process.size - Metaspace - Overhead |
taskmanager.memory.task.heap.size | — | Task Heap | 用户代码内存,无 States 时默认 ~512MB |
taskmanager.memory.managed.size | — | Managed Memory | RocksDB 场景建议占 flink.size 的 40-50% |
taskmanager.memory.managed.fraction | 0.4 | Managed Memory 占比 | RocksDB → 0.4~0.5, HashMap → 0.0 |
taskmanager.memory.network.fraction | 0.1 | Network Memory 占比 | 高吞吐场景可增大到 0.15-0.2 |
taskmanager.memory.network.min | 64MB | Network Memory 最小值 | |
taskmanager.memory.network.max | 1GB | Network Memory 最大值 | |
taskmanager.memory.jvm-overhead.fraction | 0.1 | JVM Overhead 占比 | 线程多时增大到 0.15-0.2 |
taskmanager.memory.jvm-overhead.min | 192MB | JVM Overhead 最小值 | |
taskmanager.memory.jvm-overhead.max | 1GB | JVM Overhead 最大值 | |
taskmanager.memory.framework.heap.size | 128MB | Framework Heap | |
taskmanager.memory.framework.off-heap.size | 128MB | Framework Off-Heap | |
taskmanager.memory.metaspace.size | 256MB | Metaspace | 大量 UDF 时可增大 |
7. 源码导航(完整版)
| 文件 | 关键类/方法 | 职责 |
|---|---|---|
flink-runtime/.../memory/MemoryManager.java | allocatePages(), reserveMemory(), getSharedMemoryResourceForManagedMemory() | Managed Memory 管理 |
flink-runtime/.../memory/MemorySegment.java | allocateOffHeapUnsafeMemory(), free() | 内存段 |
flink-runtime/.../memory/TaskExecutorMemoryConfiguration.java | — | TM 内存配置计算 |
flink-runtime/.../memory/UnsafeMemoryBudget.java | reserve(), release(), getAvailableMemorySize() | 内存预算 |
flink-runtime/.../network/buffer/NetworkBuffer.java | — | 网络缓冲 |
flink-runtime/.../network/buffer/NetworkBufferPool.java | createBufferPool(), requestMemorySegment() | 全局网络缓冲池 |
flink-runtime/.../network/buffer/LocalBufferPool.java | requestMemorySegment(), recycle() | 本地缓冲池 |
flink-runtime/.../network/buffer/BufferPool.java | getAvailableFuture() | 反压信号 |
flink-state-backends/flink-statebackend-rocksdb/.../RocksDBOptionsFactory.java | createDBOptions(), 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 的聚合结果)。
下一步
- 资源调度: doc-11-resource-scheduling.md — Slot 管理、K8s/YARN 模式
- 状态恢复: doc-13-state-recovery-and-rescale.md