我们在数据湖中引入了Apache Iceberg,本以为能解决数据管理的问题,结果踩了很多坑。本文记录了我们在使用Iceberg过程中遇到的各种问题,包括小文件问题、Schema演进、并发写入、性能问题等,以及每个问题的排查过程和解决方案。如果你也在用或者打算用Iceberg,希望这些踩坑经验能帮你少走弯路。

一、为什么用Iceberg

先说说我们为什么引入Iceberg。

我们的数据湖最开始是用Hive表管理的,用了几年之后问题越来越多。Hive表不支持ACID事务,写入的时候会有脏读。Schema修改很麻烦,加个字段要改很多地方。不支持时间旅行,想查历史数据很困难。小文件问题严重,查询越来越慢。

后来我们了解到了Apache Iceberg,它是一个开源的表格式,专门为数据湖设计的。它支持ACID事务、Schema演进、时间旅行、隐藏分区、小文件合并等特性,正好能解决我们的问题。于是我们决定把核心的几张表迁移到Iceberg。

迁移的过程还算顺利,但是上线之后踩了很多坑,有好几次熬夜到凌晨才解决。下面就来分享这些踩坑经历。

二、坑1:小文件爆炸

第一个大坑是小文件问题。

问题现象

上线一周之后,我们发现查询速度越来越慢。本来几秒就能查完的查询,后来要几分钟。查看文件系统,发现一个分区下面有上万个小文件,每个文件只有几KB或者几十KB。

原因分析

我们的写入任务是每5分钟运行一次的流式写入。每次写入都会生成新的数据文件,但是不会自动合并。时间一长,小文件就越来越多。查询的时候要打开上万个文件,每个文件都有打开和读取的开销,所以查询很慢。

解决方案

我们用了Iceberg的rewritedatafiles操作来合并小文件。写了一个定时任务,每天凌晨对小文件多的表执行合并操作,把小文件合并成128MB左右的大文件。

CALL system.rewrite_data_files(
  table => 'db.table',
  options => map('target-file-size-bytes', '134217728')
)

同时,我们调整了流式写入的策略,把写入间隔从5分钟改成了30分钟,减少文件生成的频率。还开启了Iceberg的自动提交合并,在写入的时候自动合并小文件。

效果

合并之后,文件数量减少了90%,查询速度提升了好几倍。这个问题算是解决了,但是要持续监控小文件数量,定期合并。

三、坑2:Schema演进的兼容性问题

第二个坑是Schema演进的问题。

问题现象

我们给一张表加了一个字段,结果历史数据的查询报错了,说找不到这个字段。而且有些老的Spark任务也跑不起来了。

原因分析

Iceberg虽然支持Schema演进,但是有一些限制。比如不能修改字段的类型(除非是安全的类型提升,比如int改bigint),不能删除字段(只能重命名为隐藏字段),字段的ID是不变的。

我们的问题是,加字段的时候用了Spark的ALTER TABLE ADD COLUMN,但是有些老的读取端用的是旧版本的Spark,不支持Iceberg的Schema演进,所以读的时候报错了。

解决方案

我们升级了所有读取端的Spark和Iceberg版本,确保都支持Schema演进。同时制定了Schema变更的规范:

  • 加字段必须设置默认值,这样历史数据查询的时候会用默认值填充
  • 不允许修改字段类型
  • 不允许删除字段,不需要的字段重命名为deprecatedxxx
  • Schema变更之前要通知所有的数据使用方

效果

规范了Schema变更流程之后,再也没有出现过兼容性问题。Iceberg的Schema演进确实很方便,但是要确保所有读取端都支持。

四、坑3:并发写入冲突

第三个坑是并发写入的问题。

问题现象

有两个写入任务同时写同一张表,结果其中一个任务失败了,报错说"Commit failed: conflict detected"。

原因分析

Iceberg支持ACID事务,但是并发写入的时候会有冲突检测。如果两个写入任务同时修改同一个分区或者同一个文件,后提交的那个会失败。

我们的两个写入任务,一个是实时数据写入,一个是批量数据补写,它们写的是同一张表的同一个分区,所以发生了冲突。

解决方案

我们调整了写入策略,把实时写入和批量补写错开时间,不要同时运行。同时,我们给批量补写加了重试机制,如果提交失败,等待几秒后重试。

对于必须并发写入的场景,我们用了Iceberg的行级删除(row-level delete)来避免冲突。或者把不同来源的数据写到不同的表,查询的时候再合并。

效果

错开写入时间之后,并发冲突的问题基本解决了。偶尔还是会有冲突,但是重试机制能处理。

五、坑4:分区字段选错了

第四个坑是分区字段的问题。

问题现象

有一张表查询特别慢,不管怎么优化都不行。后来发现是分区字段选错了。

原因分析

我们最开始按照日期分区,每天一个分区。但是这张表的数据量不大,每天只有几万条数据,每个分区只有一个小文件。查询的时候虽然用了分区过滤,但是因为分区太多,元数据管理的开销很大。

而且我们的查询经常是按用户ID查的,不按日期查,所以分区过滤没有起到作用,每次都要扫描所有分区。

解决方案

我们重新设计了分区策略。对于数据量小的表,不分区或者按月份分区。对于经常按用户ID查询的表,用了Iceberg的隐藏分区(hidden partitioning),按用户ID的哈希值分区,这样按用户ID查询的时候能命中分区。

CREATE TABLE db.table (
  user_id bigint,
  dt date,
  data string
)
PARTITIONED BY (bucket(16, user_id), months(dt))

效果

重新分区之后,查询速度提升了很多。分区设计是数据湖表设计的关键,一定要根据查询模式来选择分区字段,不要想当然地按日期分区。

六、坑5:元数据膨胀

第五个坑是元数据膨胀的问题。

问题现象

运行了几个月之后,我们发现Iceberg表的元数据目录越来越大,有几个表的元数据都有几个GB了。查询的时候加载元数据很慢,影响查询性能。

原因分析

Iceberg每次提交都会生成新的元数据文件(snapshot、manifest等),而且默认不会自动删除旧的元数据。时间一长,元数据就越来越多。

虽然Iceberg有元数据过期的机制,但是我们没有配置,所以旧的元数据一直保留着。

解决方案

我们配置了元数据过期策略,保留最近7天的快照,更早的自动删除。

CALL system.expire_snapshots(
  table => 'db.table',
  older_than => TIMESTAMP '2021-04-01 00:00:00',
  retain_last => 10
)

写了一个定时任务,每天凌晨对所有Iceberg表执行元数据过期操作。同时,我们也调整了manifest文件的合并策略,减少manifest文件的数量。

效果

元数据过期之后,元数据的大小减少了80%,查询时加载元数据的速度也快了很多。

七、坑6:Spark版本兼容性

第六个坑是Spark版本兼容性的问题。

问题现象

我们升级了Spark版本之后,Iceberg表读写都报错了,说类找不到或者方法不存在。

原因分析

Iceberg对Spark版本有要求,不同的Iceberg版本支持不同的Spark版本。我们升级Spark的时候没有同步升级Iceberg,导致版本不兼容。

而且Iceberg的Spark集成有两种方式:Spark Datasource V2和Spark Catalog,两种方式的配置和API不一样,混用也会出问题。

解决方案

我们查阅了Iceberg的官方文档,找到了和Spark版本匹配的Iceberg版本,然后同步升级。同时,我们统一了所有任务的Iceberg集成方式,都用Spark Catalog,避免混用。

我们还建立了版本兼容性矩阵,记录每个环境的Spark和Iceberg版本,升级之前先验证兼容性。

效果

版本统一之后,兼容性问题解决了。以后升级的时候也会先查兼容性矩阵,不会再盲目升级。

八、坑7:删除操作性能差

第七个坑是删除操作的问题。

问题现象

我们需要删除表中的一些数据(比如用户注销之后删除用户数据),用了Iceberg的DELETE FROM语句,结果删除操作非常慢,删几万条数据要几十分钟。

原因分析

Iceberg的删除操作是复制式的(copy-on-write),删除的时候要读取整个文件,过滤掉要删除的行,然后写一个新文件。如果文件很大,删除就很慢。

而且我们的删除条件不是分区字段,所以不能直接删分区,必须逐行过滤。

解决方案

对于大量删除的场景,我们改用了行级删除(position delete),先把要删除的行的位置写到删除文件里,查询的时候再过滤。这样删除操作很快,但是查询的时候会有额外的开销。

我们还优化了删除策略,把删除操作和小文件合并一起做。删除的时候顺便合并小文件,减少文件数量。

对于可以按分区删除的场景,尽量用分区删除(DROP PARTITION),速度很快。

效果

用了行级删除之后,删除操作从几十分钟降到了几秒。查询性能虽然有一点下降,但是在可接受范围内。

九、经验总结

踩了这么多坑之后,我们总结了一些使用Iceberg的经验。

1. 小文件管理是重中之重

Iceberg的小文件问题是最常见也最影响性能的问题。一定要有定期合并小文件的机制,同时控制写入频率,避免生成太多小文件。

2. 表设计要提前想清楚

分区字段、排序字段、Schema设计,这些都要在创建表之前想清楚。虽然Iceberg支持修改,但是修改的成本很高,特别是分区字段,基本上改不了。

3. 版本兼容性要注意

Iceberg、Spark、Hive、Flink这些组件的版本要匹配,升级之前一定要查兼容性文档,不要盲目升级。

4. 元数据要定期清理

Iceberg的元数据会不断增长,一定要配置元数据过期策略,定期清理旧的快照和元数据。

5. 并发写入要控制

Iceberg支持ACID,但是并发写入还是会有冲突。要合理安排写入任务,避免同时写同一张表的同一个分区。

6. 监控要到位

要监控Iceberg表的小文件数量、元数据大小、快照数量、查询性能等指标,发现问题及时处理,不要等问题严重了再去解决。

十、写在最后

Apache Iceberg是一个很优秀的数据湖表格式,它解决了Hive表的很多痛点。但是它也不是银弹,使用的时候还是会遇到各种问题。

我们踩的这些坑,大部分都是因为对Iceberg的理解不够深入,或者没有做好运维监控。只要提前了解常见的问题,做好规划和监控,Iceberg还是很好用的。

如果你也在考虑用Iceberg,建议先在测试环境充分验证,熟悉它的特性和限制,再逐步迁移到生产环境。迁移的时候先迁非核心的表,积累经验之后再迁核心表。

最后用一句话结束本文:"踩坑不可怕,可怕的是踩了坑不总结。"希望我们的踩坑经验能帮你少走弯路,用好Iceberg,建好数据湖。