引言
数据湖上线时总是很美好:写入流畅、查询飞快。三个月后,查询开始变慢、元数据读取超时、存储账单虚高。问题几乎总是同一个——没人做维护。流式写入每几分钟产出一个小文件,Upsert 留下一堆 delete 文件,快照越积越多,元数据越读越慢。
数据湖不是"写完就不用管"的存储;它更像一张需要定期保养的表——合并、清理、回收,缺一不可。
本文把数据湖运维讲透:为什么会产生小文件、怎么合并、怎么清理元数据、怎么回收存储、不同表格式的差异,以及如何把维护做成自动化的流水线。
一、数据湖为什么需要持续运维
1.1 写入模式决定维护需求
| 写入模式 | 产出特征 | 维护需求 |
|---|---|---|
| 批量追加 | 文件较大、数量可控 | 低频合并 |
| 流式写入 | 高频小文件 | 高频合并 |
| Upsert | 小文件 + delete 文件 | 合并 + 清理 delete |
| 频繁更新 | 版本膨胀 | 快照过期 + 合并 |
写入越"实时",产出的文件越小、越多,维护需求越高。这是流批一体湖仓的固有代价。
1.2 性能退化的三个信号
# [ ] 查询扫描文件数暴增, 单查询打开上千文件
# [ ] 元数据读取变慢, 计划阶段就耗时数秒
# [ ] 存储用量远超实际数据量(快照/孤儿残留)
一旦出现这些信号,说明维护已经欠账,需要尽快补上。
1.3 维护的本质
维护的本质是用后台的写放大换前台查询的稳定。合并会重写数据、消耗 IO 与算力,但换来的是文件数下降、查询变快、元数据瘦身。这是一笔必须定期偿还的账。
二、小文件问题的成因与代价
2.1 小文件从哪来
- 高频流写:Flink 每 1-5 分钟一次 checkpoint,每次都产出一批文件。
- 分区过细:按小时甚至分钟分区,单分区数据量很小。
- Upsert 的 delete 文件:每次更新生成 delete 文件,本身也是"小文件"。
- 过度并行:写入并行度远大于数据量,每个 writer 产出一个小文件。
-- 查看某表的分区文件分布, 定位小文件重灾区
SELECT partition, count(*) AS files, avg(file_size_in_bytes)/1024/1024 AS avg_mb
FROM lake.orders.files
GROUP BY partition
HAVING count(*) > 100
ORDER BY files DESC
LIMIT 20;
2.2 小文件的代价
| 代价 | 表现 | 影响 |
|---|---|---|
| 元数据膨胀 | Manifest 记录数暴增 | 计划阶段变慢 |
| 打开文件开销 | 每文件一次 IO | 查询延迟升高 |
| 列统计失效 | 文件小统计粒度细 | 剪枝收益下降 |
| 存储元数据开销 | 对象存储按请求计费 | 成本上升 |
2.3 目标文件大小
经验值:目标 128MB,下限 32MB。太小合并收益低,太大影响并行度与增量读取。
# 文件大小参考
# < 8MB 严重小文件, 优先合并
# 8-32MB 偏小, 可合并
# 32-128MB 健康区间
# > 512MB 偏大, 可能影响并行
三、合并(Compaction)策略与实现
3.1 binpack:纯物理合并
binpack 只做物理层面的文件合并,不改变数据排序,速度快、成本低,是日常维护的主力。
-- Iceberg: binpack 合并
CALL lake.system.rewrite_data_files(
table => 'ods.orders',
strategy => 'binpack',
min_file_size_in_bytes => 33554432, -- 32MB 以下视为小文件
max_file_size_in_bytes => 134217728, -- 合并到 128MB 封顶
where => "dt >= '2026-09-01'" -- 限定范围, 控制成本
);
3.2 sort:合并时重排
sort 在合并的同时按指定列排序,能显著提升后续查询的剪枝与压缩效率,但成本更高。
-- Iceberg: sort 合并, 按 event_time 排序
CALL lake.system.rewrite_data_files(
table => 'ods.orders',
strategy => 'sort',
sort_order => 'event_time DESC, user_id ASC',
min_file_size_in_bytes => 33554432
);
3.3 增量合并与全量合并
# 增量合并: 只处理最近变更的分区, 成本低, 日常跑
# 全量合并: 重写整表布局, 成本高, 定期(如每周)跑
# 策略: 增量为主, 全量兜底历史遗留布局
3.4 合并的调度节奏
| 表类型 | 合并频率 | 策略 |
|---|---|---|
| 高频流写 | 每小时 | binpack |
| 每日批量 | 每日 | sort |
| 低频维表 | 每周 | binpack + sort |
| 历史分区 | 一次性 | 全量 sort |
关键是合并速度必须追上小文件产生速度,否则欠账越积越多。
四、快照与元数据维护
4.1 快照膨胀
每次提交都生成新快照,快照保留着历史文件清单。若不清理,历史数据文件永远不会被删除,存储持续增长,元数据文件也越来越大。
-- 查看快照数量与时间跨度
SELECT snapshot_id, committed_at, operation
FROM lake.orders.snapshots
ORDER BY committed_at DESC
LIMIT 20;
4.2 快照过期
-- 清理 7 天前、且保留最近 10 个快照
CALL lake.system.expire_snapshots(
table => 'ods.orders',
older_than => TIMESTAMP '2026-09-26 00:00:00',
retain_last => 10
);
注意:过期不可逆。保留窗口要覆盖审计与回溯需求,常见 7-30 天。
4.3 元数据文件与 Manifest 清理
快照过期后,旧的 metadata 文件与 Manifest 可能成为孤儿,需要一并清理。Iceberg 提供 rewrite_manifests 合并碎片化 Manifest,减小元数据读取开销。
-- 合并小 Manifest, 减小元数据体积
CALL lake.system.rewrite_manifests(table => 'ods.orders');
-- 清理孤儿元数据文件
CALL lake.system.remove_orphan_files(
table => 'ods.orders',
older_than => TIMESTAMP '2026-09-26 00:00:00'
);
五、孤儿文件与存储回收
5.1 什么是孤儿文件
孤儿文件是"存在于对象存储、但没有任何快照引用"的文件。来源包括:写入失败残留、被取消的任务、过期快照遗留的数据文件。
5.2 回收流程
# 回收三步
# 1. expire_snapshots: 让旧快照失效
# 2. remove_orphan_files: 删除无引用文件
# 3. 校验存储用量: 确认回收生效
# 关键: older_than 必须足够久, 避免误删正在写入的文件
-- 孤儿文件清理: older_than 设为 3 天前, 留足安全窗口
CALL lake.system.remove_orphan_files(
table => 'ods.orders',
older_than => TIMESTAMP '2026-09-30 00:00:00',
dry_run => false
);
5.3 误删风险与防护
- older_than 太近:会删掉正在写入、尚未提交的文件,导致数据损坏。
- 未停写就清理:并发写入期间清理风险极高,应在低峰或暂停写入时执行。
- 先 dry_run:生产环境先 dry-run 看清单,再真正执行。
六、Iceberg/Hudi/Delta 的运维差异
6.1 三种表格式对比
| 维度 | Iceberg | Hudi | Delta Lake |
|---|---|---|---|
| 合并命令 | rewrite_data_files | compaction | OPTIMIZE |
| 清理快照 | expire_snapshots | clean | VACUUM |
| 孤儿清理 | remove_orphan_files | clean | VACUUM |
| 元数据 | Manifest | Timeline | Log |
| 小文件治理 | binpack/sort | inline/clustering | OPTIMIZE + Z-ORDER |
6.2 Hudi 的 inline 与异步
Hudi 支持 inline compaction(写入时合并)与异步 compaction。inline 简单但拖慢写入;异步把合并放到独立任务,写入更快,但需额外运维。
-- Hudi 异步 compaction 关键配置
hoodie.compact.inline=false
hoodie.compact.schedule.inline=true
hoodie.compaction.strategy=org.apache.hudi.DefaultCompactionStrategy
hoodie.compaction.target.io=102400 -- 单次合并目标字节
6.3 Delta 的 OPTIMIZE 与 Z-ORDER
-- Delta: 合并 + Z-ORDER 多维聚簇
OPTIMIZE delta.`s3://lake/orders` ZORDER BY (user_id, dt);
VACUUM delta.`s3://lake/orders` RETAIN 168 HOURS; -- 保留 7 天
Z-ORDER 让多个查询维度都获得数据聚簇,代价是合并成本更高。
七、自动化运维流水线
7.1 用编排工具串联
维护任务应纳入调度系统,形成固定流水线,而不是靠人手动跑。
# airflow 维护 DAG 片段(抽象)
maintenance_orders:
schedule: "0 * * * *" # 每小时
tasks:
- compact_binpack # 先合并小文件
- rewrite_manifests # 再整理元数据
- expire_snapshots # 清理过期快照
daily:
- remove_orphan_files # 每日回收孤儿
- compact_sort # 每日重排
7.2 表级参数自动维护
Iceberg 支持表级维护参数,让引擎自动触发小规模合并,减少人工干预。
-- 开启写入时的自动小文件合并
ALTER TABLE ods.orders SET TBLPROPERTIES (
'write.target-file-size-bytes' = '134217728',
'commit.manifest-merge.enabled' = 'true',
'commit.manifest.min-count-to-merge' = '100'
);
7.3 分优先级调度
# 优先级策略
# P0 核心表: 每小时合并, 每日清理
# P1 一般表: 每日合并, 每周清理
# P2 归档表: 每周合并, 每月清理
# 资源紧张时优先保障 P0
八、监控指标与踩坑清单
8.1 必看指标
# [ ] 平均文件大小: 低于 32MB 触发告警
# [ ] 文件总数趋势: 持续上升说明合并没跟上
# [ ] 快照数量: 超过阈值(如 100)需清理
# [ ] 元数据文件大小: 增长过快需 rewrite_manifests
# [ ] 孤儿文件数: 定期扫描对象存储核对
# [ ] 合并任务耗时与成功率
8.2 生产踩坑
- 合并任务与写入冲突:并发提交导致冲突重试,应错峰或限定分区。
- older_than 设太近:误删活跃文件,务必留足安全窗口。
- 只合并不清理快照:合并产出的旧文件仍被快照引用,存储不降反升。
- 全表全量合并:一次重写整表,成本失控,应限定范围。
- 忽略 delete 文件:Upsert 表只合数据不合 delete,读取仍慢。
- 维护任务无监控:任务静默失败,欠账无人知。
总结
| 维护动作 | 解决什么 | 频率 |
|---|---|---|
| binpack 合并 | 小文件过多 | 高频 |
| sort 合并 | 布局混乱 | 定期 |
| rewrite_manifests | 元数据碎片 | 定期 |
| expire_snapshots | 快照膨胀 | 定期 |
| remove_orphan_files | 存储虚高 | 定期 |
| 监控告警 | 欠账预警 | 持续 |
数据湖运维的核心认知是:维护不是一次性的,而是持续运营。小文件会随写入不断产生,快照会随提交不断累积,只有把合并、清理、回收做成自动化的固定流水线,并配上监控告警,数据湖才能长期保持查询性能与成本的可控。理解每种表格式的运维命令差异,是落地这套流水线的前提。
参考与延伸阅读
- Apache Iceberg 官方文档:维护存储过程与表配置
- Apache Hudi 官方文档:Compaction 与 Clean 策略
- Delta Lake 官方文档:OPTIMIZE 与 VACUUM
- Apache Iceberg 深度剖析 — 元数据模型基础
- 数据湖技术 — 表格式与存储选型
- 数据 FinOps 成本优化 — 存储成本治理
- 数据可观测性 — 维护任务监控方法
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。