一、问题现象

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 怎么办?

按优先级排查:

  1. 减少 slot 数numberOfTaskSlots: 64),降低并发 ParquetWriter 数量
  2. 减小 Row Group 大小(ChunJun JSON 配置中设置 rowGroupSize),降低单次 Snappy 压缩的 Direct Memory 峰值
  3. 换压缩算法(SNAPPY → GZIP),GZIP 使用堆内压缩,不消耗 Direct Memory
  4. 检查是否 Direct Memory 泄漏(OOM 频率随时间递增),可能是 Parquet Writer 异常后未关闭

七、经验总结

  1. Flink 的 task.off-heap 默认值为 0,任何使用 Direct ByteBuffer 的用户代码(Parquet Snappy、Kryo、Netty 自定义等)都必须显式配置此值

  2. Flink 自动设置 MaxDirectMemorySize,公式 = task.off-heap + framework.off-heap + network。不配 task.off-heap 就等于没给用户代码留 Direct Memory

  3. Direct Memory OOM 跟任务运行时长无关,跟数据吞吐量有关。低吞吐时 GC 能及时回收 Direct ByteBuffer,高吞吐时分配速度超过回收速度

  4. ChunJun 的 Parquet + Snappy 写入是 Direct Memory OOM 的高危场景:Snappy 使用 Direct ByteBuffer 压缩、HiveOutputFormat 多分区复用成倍放大、ChunJun 不暴露 Parquet 内存池配置

  5. HeapDumpOnOutOfMemoryError 务必开启,虽然 heap dump 不包含 Direct Memory 内容,但能看到持有 Direct ByteBuffer 引用的对象,辅助定位


环境:Flink 1.16.1 / ChunJun(dtstack)/ Parquet 1.11.1 / JDK 1.8.0_301 / CDH 3 节点集群

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐