Flink standalone + ChunJun 写入 HDFS 之 Direct Buffer Memory OOM 排查与调优
一、问题现象
Flink 1.16.1 集群(standalone,3 节点)运行 ChunJun 同步任务,向 HDFS 写入 Parquet 格式文件时,TaskManager 报 OOM 并崩溃:
java.lang.OutOfMemoryError: Direct buffer memory
at java.nio.Bits.reserveMemory(Bits.java:695)
at java.nio.ByteBuffer.allocateDirect(ByteBuffer.java:311)
at org.apache.parquet.hadoop.codec.SnappyCompressor.setInput(SnappyCompressor.java:97)
at org.apache.parquet.hadoop.CodecFactory$HeapBytesCompressor.compress(CodecFactory.java:165)
at org.apache.parquet.column.impl.ColumnWriterV1.flush(ColumnWriterV1.java:238)
at org.apache.parquet.hadoop.InternalParquetRecordWriter.flushRowGroupToStore(InternalParquetRecordWriter.java:167)
...
at com.dtstack.chunjun.connector.hdfs.sink.HdfsParquetOutputFormat.writeSingleRecordToFile(HdfsParquetOutputFormat.java:192)
诡异之处: 同样的 Parquet + Snappy 写入方式,有的任务能跑一整晚没问题,有的任务几十分钟就 OOM。
二、根因定位
2.1 调用链分析
从堆栈看,OOM 发生在 Parquet 的 Snappy 压缩环节:
ChunJun writeRecord()
→ ParquetWriter.write() // 写入一条记录
→ ColumnWriterV1.flush() // Row Group 满 128MB,触发 flush
→ SnappyCompressor.setInput()
→ ByteBuffer.allocateDirect() // ← 这里分配 Direct Memory
关键点: 这是用户代码路径上的 Direct Memory 分配,不是 Flink 框架管理的网络缓冲区。
2.2 为什么默认配置一定会 OOM
查看原始配置,发现 taskmanager.memory.task.off-heap.size 没有配置,默认值为 0。
Flink 启动 TM 时会自动计算并设置 JVM 参数 -XX:MaxDirectMemorySize:
MaxDirectMemorySize = task.off-heap + framework.off-heap + network
= 0 + 128m + 2560m
= 2688m
而这三部分的实际占用:
网络缓冲区(Flink 启动后逐步分配) -2560m
框架 off-heap - 128m
─────────────────────────────────────────
剩余给 Parquet/Snappy 的 Direct Memory = 0m
MaxDirectMemorySize 总共 2688m,被网络缓冲区和框架占完,用户代码一个字节的空间都没有。
三、为什么有的任务能跑一晚上
核心原因:数据吞吐量不同。
3.1 Parquet 写入的 Direct Memory 是"脉冲式"使用
正常写入(堆内存): record → ColumnWriter 内存缓冲(heap)→ 累积到 128MB
↑ 不占 Direct Memory
Row Group Flush 瞬间: Snappy 压缩 → allocateDirect() → Direct Memory 短暂飙升
↑ 毫秒~秒级
Flush 结束: Direct ByteBuffer 引用释放 → 等 GC 回收 → Direct Memory 归零
3.2 低吞吐任务(能跑一晚上)
数据量小 → Row Group flush 每 10 分钟一次
→ 6 个 slot 不同时 flush
→ 每次峰值 ~100MB
→ GC 来得及回收,不累积
→ 总占用 ≈ 2560(网络) + 128(框架) + 100(Snappy) = 2788m ≈ 2688m
→ 勉强没超限,但已经在悬崖边上
3.3 高吞吐任务(快速 OOM)
数据量大 → Row Group flush 每 30 秒一次
→ 6 个 slot 接近同时 flush
→ 峰值 ~600MB+
→ GC 来不及回收,累积超限
→ 总占用 ≈ 2560 + 128 + 600 = 3288m > 2688m
→ OOM!
3.4 总结
| 因素 | 不容易 OOM | 容易 OOM |
|---|---|---|
| 数据吞吐量 | 低(慢速源) | 高(高速源) |
| 并发 subtask 数 | 少 | 多(6 slot 满载) |
| Row Group 大小 | 小 | 大(128MB/256MB) |
| 写入分区数 | 单分区 | 多分区(HiveOutputFormat 多路复用) |
数据量越大的任务 OOM 越快,数据量小的任务能撑很久,但本质上都在悬崖边上——只是低吞吐时刚好没踩到。
四、ChunJun 代码分析(影响 Direct Memory 的关键特征)
通过阅读 ChunJun 源码,发现以下特征直接影响 Direct Memory 配置:
4.1 Parquet Writer 创建(无 Direct Memory 管理)
// HdfsParquetOutputFormat.nextBlock()
ExampleParquetWriter.builder(writePath)
.withCompressionCodec(compressionCodecName) // 默认 SNAPPY
.withRowGroupSize(hdfsConfig.getRowGroupSize()) // 默认 128MB
...
- ChunJun 没有任何 Direct Memory 管理、释放或限制的代码
- 所有 Direct Memory 完全依赖 Parquet 库内部和 JVM 的 MaxDirectMemorySize
4.2 HiveOutputFormat 多路复用
// HiveOutputFormat 维护多个 BaseHdfsOutputFormat 实例
Map<String, Pair<String, BaseHdfsOutputFormat>> outputFormatMap;
一个 HiveOutputFormat 按表/分区持有多个独立的 ParquetWriter,每个 Writer 独立分配 Direct Memory。写入多分区时,Direct Memory 消耗成倍增长。
4.3 写异常时 Writer 不关闭
// HdfsParquetOutputFormat.writeSingleRecordToFile()
writer.write(group); // 如果这里抛 IOException
// writer 不会被关闭,其 Direct Memory 持续占用
// 直到下一个 checkpoint flush 或任务结束
4.4 无 Parquet 内存池配置
ChunJun 不暴露 parquet.memory.max.allocation 等参数,用户无法在 ChunJun 层面控制 Parquet 的内存使用。
4.5 默认压缩为 Snappy
// CompressType 枚举:compress 为空时,Parquet 默认选 SNAPPY
Snappy 是唯一使用 Direct Memory 的压缩算法。GZIP/BZIP2 使用堆内压缩,不会触发此问题。
五、解决方案
5.1 核心配置调整
修改所有节点的 flink-conf.yaml:
# 1. 新增:用户代码堆外内存(给 Parquet/Snappy 的 Direct Memory 预算)
# 128MB row group: 设为 2048m
# 256MB row group: 设为 4096m
taskmanager.memory.task.off-heap.size: 4096m
# 2. 减少 task.heap,为 task.off-heap 腾空间
# 从 20480m 减至 16000m
taskmanager.memory.task.heap.size: 16000m
# 3. 适当增大 framework.off-heap(高并行度场景)
taskmanager.memory.framework.off-heap.size: 256m
# 4. JVM 参数:显式设置 MaxDirectMemorySize
# = task.off-heap(4096) + framework.off-heap(256) + network(2560)
env.java.opts.taskmanager: >-
-XX:+UseG1GC
-XX:MaxGCPauseMillis=200
-XX:InitiatingHeapOccupancyPercent=40
-XX:MaxDirectMemorySize=6912m
-XX:+HeapDumpOnOutOfMemoryError
-XX:HeapDumpPath=/data/flink1.16.1/log/
5.2 内存分配验证
32GB TM 进程内存分配:
task.heap = 16000m 用户代码堆内存
task.off-heap = 4096m Parquet/Snappy Direct Memory
managed = 8192m RocksDB 状态后端 + 排序
framework.heap = 128m Flink 框架
framework.off-heap = 256m Flink 框架 Direct Memory
network = 2560m 网络缓冲区
JVM overhead = 1536m JVM 自身开销
合计 = 32768m ✓
Direct Memory 总预算 = 4096 + 256 + 2560 = 6912m
5.3 配置对照表(按 Row Group 大小)
| 配置项 | 128MB Row Group | 256MB Row Group |
|---|---|---|
| task.off-heap | 2048m | 4096m |
| task.heap | 18048m | 16000m |
| framework.off-heap | 256m | 256m |
| MaxDirectMemorySize | 4864m | 6912m |
六、Direct Memory 计算 Q&A
Q1: task.off-heap 是所有 slot 共享还是每个 slot 独占?
共享。 Flink 没有 slot 级别的 Direct Memory 隔离。task.off-heap = 4096m 是整个 TM 进程中所有 slot 的 Direct Memory 总预算。
┌─────────────────── TM 进程(1 个 JVM)──────────────────┐
│ │
│ task.off-heap = 4096m(所有 slot 共享,无隔离) │
│ │
│ Slot-1 ──┐ │
│ Slot-2 ──┤ │
│ Slot-3 ──┼── 共同从 4096m 中分配 Direct ByteBuffer │
│ Slot-4 ──┤ │
│ Slot-5 ──┤ │
│ Slot-6 ──┘ │
└──────────────────────────────────────────────────────────┘
Q2: 为什么 OOM 的 heap dump 看不到 Direct Memory 内容?
Direct ByteBuffer 的数据存储在 JVM 堆外(native memory),heap dump 只能看到 Java 堆中的对象引用。但 heap dump 能看到哪些对象持有了 Direct ByteBuffer 引用,从而定位是哪些 ParquetWriter 占用了 Direct Memory。
Q3: 如果调大后仍然 OOM 怎么办?
按优先级排查:
- 减少 slot 数(
numberOfTaskSlots: 6→4),降低并发 ParquetWriter 数量 - 减小 Row Group 大小(ChunJun JSON 配置中设置
rowGroupSize),降低单次 Snappy 压缩的 Direct Memory 峰值 - 换压缩算法(SNAPPY → GZIP),GZIP 使用堆内压缩,不消耗 Direct Memory
- 检查是否 Direct Memory 泄漏(OOM 频率随时间递增),可能是 Parquet Writer 异常后未关闭
七、经验总结
-
Flink 的
task.off-heap默认值为 0,任何使用 Direct ByteBuffer 的用户代码(Parquet Snappy、Kryo、Netty 自定义等)都必须显式配置此值 -
Flink 自动设置
MaxDirectMemorySize,公式 = task.off-heap + framework.off-heap + network。不配 task.off-heap 就等于没给用户代码留 Direct Memory -
Direct Memory OOM 跟任务运行时长无关,跟数据吞吐量有关。低吞吐时 GC 能及时回收 Direct ByteBuffer,高吞吐时分配速度超过回收速度
-
ChunJun 的 Parquet + Snappy 写入是 Direct Memory OOM 的高危场景:Snappy 使用 Direct ByteBuffer 压缩、HiveOutputFormat 多分区复用成倍放大、ChunJun 不暴露 Parquet 内存池配置
-
HeapDumpOnOutOfMemoryError务必开启,虽然 heap dump 不包含 Direct Memory 内容,但能看到持有 Direct ByteBuffer 引用的对象,辅助定位
环境:Flink 1.16.1 / ChunJun(dtstack)/ Parquet 1.11.1 / JDK 1.8.0_301 / CDH 3 节点集群
更多推荐




所有评论(0)