Flink 滑动窗口作业 TaskManager RocksDB 状态膨胀 Bug 分析
- 文档版本:v1.0
- 编写日期:2026-07-21
- 适用范围:基于 Flink 1.19 的滑动窗口聚合作业(RocksDB State Backend)
- 文档定位:技术根因分析 + 修复方案
1. 问题现象
1.1 最显性表现:Flink Checkpoint 体积持续增长
运维 dashboard 上最先观察到的是——该作业的 checkpoint 大小长时间不停地涨,没有收敛迹象:
- 启动初期 checkpoint 体积正常(MB~GB 级别)
- 运行 N 小时后,checkpoint 体积持续单调上升,明显与业务规模脱钩
- 即便输入流量稳定,checkpoint 仍持续增长——这是典型的状态后端内部 churn,不是数据规模问题
这一条是最早、最易发现的信号,也是触发本次排查的直接原因。
1.2 下钻后的内层证据
顺着 checkpoint 体积异常这条线,追到 TaskManager 本地 RocksDB 目录,发现:
- 实际数据 sst 文件总数很少,4 个 sst 文件合计约64KB(26K + 6K + 1.3K + 31K)
- MANIFEST 单文件达到 447MB,且持续单调增长
- 目录归属路径示例:
即某 WindowOperator 的第 3 个并行度、第 3 个 subtask 的 RocksDB 实例/tmp/tm_xxx/job_xxx_op_WindowOperator_xxx/__3_3__uuid_xxx/db/
1.3 影响
- 磁盘占用失控:按约 460MB/天的增速,单 subtask 一周可再涨 3GB+,多算子多并行度叠加后非常可观
- Checkpoint 自身受损:体积越来越大 → barrier 对齐耗时上升 → checkpoint 时延拉长 → 长时间运行后 checkpoint 容易超时失败
- 故障恢复时间失控:RocksDB 启动需要重放 MANIFEST,MANIFEST 越大恢复越慢
- 磁盘写满风险:极端情况下直接导致 TaskManager 异常退出
1.4 出现规律
- 与作业运行时间正相关,越往后越严重
- 不依赖 checkpoint 行为即可复现(即便关掉 checkpoint,本地 MANIFEST 仍在涨——只是这次不会再被上传到外部存储放大成 checkpoint 体积)
2. 排查过程
按时间顺序记录本次问题的调查路径,便于复盘与他人复用排查套路。
2.1 第一步:从运维侧看到 checkpoint 异常
在 Flink Web UI / 监控大盘上观察到:
- 该作业checkpoint 大小持续单调上涨,没有收敛
- 业务输入流量、Key 数量在同一时期并未显著变化
- 多个 WindowOperator 都有类似趋势,但并非所有算子均匀分布——其中 WindowOperator 的某些并行度/subtask 增长尤为明显
判断:状态后端内部在持续"折腾"——纯数据规模增长不会跟输入解耦成这样,更像是状态后端元数据层在做大量重写。
2.2 第二步:下钻到 TaskManager 本地磁盘
顺着 checkpoint 异常选了一个增长最严重的 subtask,上到对应 TaskManager:
ls-lh/tmp/tm_xxx/job_xxx_op_WindowOperator_xxx/__3_3__uuid_xxx/db/观察到反常现象:
MANIFEST*单个文件447MB*.sst数据文件仅 4 个,合计64KB
这种"sst 几乎为空、MANIFEST 巨大"的组合——数据本身不大,问题在 RocksDB 的元数据/版本变更日志,而不是用户数据。
2.3 第三步:grep RocksDB LOG 验证 churn 频率
RocksDB 会在 INFO 级 LOG 中记录Flush/Compaction事件。直接在出问题实例上 grep:
grep-c"processing_window-timers"LOG LOG.old.*grep-c"Moving #"LOG LOG.old.*grep-cE"Flushing memtable|Compacting"LOG LOG.old.*得到关键数字(仅 LOG 当前保留的窗口,约 8.5 分钟):
| 关键字 | 出现次数 |
|---|---|
Flushing memtable|Compacting | 30,598 |
processing_window-timers | 22,500 |
Moving # | 0 |
判断:
- Flush + Compaction 速率 ≈60 次/秒—— 远超正常水平,意味着 RocksDB 在疯狂写盘
- 大量 timer 写入与作业的窗口特征一致
- 与此同时"Moving #" = 0,提示 compaction 走的是真实合并而非 trivial move
2.4 第四步:交叉验证与归因
把已知因素排在一起对照:
| 现象 | 候选假设 1:数据规模大 | 候选假设 2:写入 churn 元数据层 | 候选假设 3:纯内存不足 |
|---|---|---|---|
| sst 总大小 64KB | ❌ 应很大 | ✅ churn 元数据不影响 sst | ✅ sst 由真实数据量决定 |
| MANIFEST 447MB | ❌ 不应远大于 sst | ✅ 每次 flush/compact 都加 VersionEdit | ✅ 写入频率高就涨 |
| Flush 速率 60/s | ❌ 正常流量下不该这么高 | ✅ 是结果,不是因 | ✅ 内存不足 → memtable 频繁刷出 |
| 业务输入未变 | — | — | ✅ 与本次内存变更(12G→8G)吻合 |
三条假设里,只有"内存不足"能同时解释所有现象,且与近期变更记录一致。同时processing_window-timers频次高 → 提示 sliding window + CountTrigger 共同放大了写入量,属于助推器。
2.5 第五步:定位真因
结合:
- 近期变更记录:taskmanager 物理内存由 12G 调小到 8G
- 算子特征:三套滑动窗口(10/30/60min,slide=30s)+ CountTrigger.of(1) + 最后合并使用 ProcessWindowFunction
锁定根因链:内存不足 → memtable flush 风暴 → MANIFEST 膨胀 → checkpoint 上涨。
3. 根因分析
3.1 真正根因:TaskManager 物理内存不足,memtable 被持续刷出
核心机制链:
TaskManager 物理内存由 12G → 8G(缩减) ↓ 实际写入量未减少,但可用堆外内存(RocksDB write buffer / block cache)紧张 ↓ RocksDB memtable 频繁写满,触发 flush 到 L0 ↓ L0 小文件堆积速度 > compaction 消化速度 ↓ 大量新 sst 文件被创建,MANIFEST 中产生海量 VersionEdit(增/删文件记录) ↓ "活着的数据很少(64KB)但 MANIFEST 巨大(447MB)"的反常现象关键反直觉点:磁盘上"真正存活的用户数据"几乎没有,膨胀的不是数据本身,而是版本变更日志——RocksDB 每次 flush / compaction 都会在 MANIFEST 中追加一条 VersionEdit 记录,频繁的 memtable flush 把 MANIFEST 写爆。
3.2 助推器(不是根因,但显著放大症状)
3.2.1 CountTrigger.of(1) 导致 200 倍输出放大
作业内使用了三套SlidingProcessingTimeWindows(10min / 30min / 60min,slide 均为 30s),配合CountTrigger.of(1):
- sliding window 语义:一条输入会被同时分配到 size/slide 个窗格
- 三套窗口合计每条输入会被复制到20 + 60 + 120 = 200 个窗格
CountTrigger.of(1)是每收到 1 条新元素就立即 fire(注意是 FIRE 不是 FIRE_AND_PURGE,reduce 累加器的 state 不会被清空)- 一条原始事件经过 union 后变成 200 条下游记录,直接推高整个流水线的写入频率
3.3 写入频率量化(基于真实日志统计)
TaskManager 上对 RocksDB LOG 取样(4 个 LOG/LOG.old 文件,约 8.5 分钟跨度):
| 事件 | 出现次数(8.5 分钟内) | 折算频率 |
|---|---|---|
Flushing memtable | Compacting | 30,598 | ≈ 60 次/秒 |
processing_window-timers提及 | 22,500 | ≈ 44 次/秒 |
Moving #(trivial move) | 0 | — |
两个关键事实:
- RocksDB LOG 文件会轮转、旧的会被清理——这 8.5 分钟只是 LOG 当前保留的窗口,不代表"23 小时里只有 30K 次"。结合文件创建时间戳精确反推,累计频率可由 ~60 次/秒 × 23 小时 ≈ 505 万次解释,按单条 VersionEdit 约 90 字节估算,倒推回 447MB 的 MANIFEST 大小在合理范围内。
- 即使按这 8.5 分钟的窗口做下限估算,每秒 60 次 memtable flush / compaction 本身已经是严重资源紧缺的典型信号——正常配置的 RocksDB 不会以这个速率运转。
3.4 行为交叉验证
- ✅ “sst 文件小(KB 级)× MANIFEST 大(百 MB 级)”——典型 churn 模式,非数据规模问题
- ✅
Flush / Compaction频次和 MANIFEST 增速匹配——次数堆量解释一切 - ✅ 列族名
processing_window-timers大量出现——窗口密集创建/清理,与 sliding window + CountTrigger 行为一致
4. 修复方案
4.1 紧急 / 短期缓解
- 恢复 taskmanager 内存:12G → 8G 是反向操作,立刻改回≥ 12G,并视写入量上调
- 重启 TaskManager:清掉现有 RocksDB 实例,从最近一次 checkpoint 重启后 MANIFEST 立刻归零——这是能立刻见效的临时手段,但根因未根治会再次涨上来
5. 复盘与防患
5.1 变更管理问题
- 建议:内存/并发/分区等容量参数进入 CI/CD 卡口,必须配套附上"本次变更对应的负载评估/压测结论"
- 建议:在 Flink UI / Prometheus 上加监控告警:
- RocksDB
Flush / Compaction速率持续 > X/s - RocksDB
WriteStall/Stalling writes because we have N level-0 files出现频次 - TM 堆外内存使用率 /
write-buffer使用率
- RocksDB
5.2 可观测性补丁
部署以下观测项可在下次提前发现同类问题:
| 指标 | 阈值建议 | 含义 |
|---|---|---|
RocksDBnum files at level | L0 > 10 持续 | 写入压力超 compaction 能力 |
RocksDBflush countrate | > 10/s 持续 | memtable 撑爆 |
| RocksDB manifest size | 单实例 > 100MB | MANIFEST 异常膨胀 |
pending compaction bytes | 持续 > 1GB | compaction backlog |
| TM native memory used / total | > 85% | 写入缓冲即将耗尽 |
5.3 后续 TODO
6. 附录:原始证据
6.1 关键文件大小(1 个 subtask 抽样)
db 目录文件清单(抽样) - 4 个 sst 文件:26K + 6K + 1.3K + 31K ≈ 64KB 合计 - MANIFEST:447MB6.2 RocksDB LOG grep 结果(仅 LOG 当前保留窗口)
grep-c"processing_window-timers"LOG LOG.old.* LOG:4493 LOG.old.1784617536347429:5846 LOG.old.1784617679552965:4952 LOG.old.1784617874869913:7209grep-c"Moving #"LOG LOG.old.*# 全部为 0grep-cE"Flushing memtable|Compacting"LOG LOG.old.* LOG:7061 LOG.old.1784617536347429:7563 LOG.old.1784617679552965:7811 LOG.old.1784617874869913:8163合计:processing_window-timers 出现 22,500 次,Flush+Compact 出现 30,598 次,时间窗口 ≈ 8.5 分钟,对应频率 ≈ 60 次/秒。