首页 / 资讯中心 / 文章详情

Apache Iceberg高并发写入一致性机制与生产实践

Apache Iceberg高并发写入一致性机制与生产实践 ★ FEATURED ARTICLE
这批数据进湖的时候同时跑着两个实时写入任务目标都是同一张Iceberg表。凌晨两点下游告警群里突然跳出一排红字某个Spark批任务提交失败报的异常是CommitFailedException后面还跟着一句“Cannot commit because the current table metadata is newer than the metadata this transaction is based on”。当时第一反应是任务并发不够、资源不够但拉完日志仔细一看根本不是资源问题而是两个任务在同一秒里抢着提交Iceberg用自己的乐观并发机制把其中一方拦了下来。这次事故让我彻底认认真真把Iceberg的并发一致性机制翻了一遍。当然理解这些机制并不只是为了排查故障更是为了在设计写入链路时知道哪些参数能调、哪些策略不能乱碰。下面把这段时间的实测和踩坑整理出来尤其是高并发场景下Iceberg是怎么保证数据一致性的以及它在生产环境里真正的边界在哪里。1. 写入即“写新不写旧”一致性机制的地基1.1 元数据三层结构一张表的“账本”是这么记的很多人刚接触Iceberg时第一个困惑是它和Hive表到底差在哪Hive表的元数据基本是一份分区目录加文件列表写一半的时候读端能看到半成品文件Iceberg不一样它给一张表维护了一套完整的元数据版本链每次变更不是去改老文件而是生成一份新的元数据描述。这套描述从上到下分四层Catalog层存的是表当前指向哪个版本的元数据文件相当于账本封面上的“最新版本号”。Table Metadata层一个JSON文件里面记录了表的schema、分区信息、快照列表以及当前快照的ID。Manifest List层一次快照对应的清单列表相当于这次快照涉及哪些Manifest文件。Manifest File层每个Manifest里记录了一批数据文件的路径、分区范围、列统计信息。数据文件本身是不可变的。无论是INSERT、DELETE还是UPDATEIceberg的策略都是“写新文件、不改老文件”所有变更靠“生成一套新的Manifest来描述新文件集合”来完成。这个设计很像快递分拣中心的账本货架上的包裹数据文件入库之后不挪动发货单Manifest不断重印每一版发货单都完整描述当前货架上应该有哪些包裹。1.2 提交瞬间到底发生了什么从指针切换看原子性只要理解了“数据不可变、元数据可版本化”这个前提再去看提交Commit就清晰多了。一次正常的INSERT写入底层动作大概是这个顺序写入任务先向文件系统写若干个新的数据文件Parquet/ORC/Avro。根据这些数据文件生成对应的Manifest文件。生成或更新Manifest List把新的Manifest合并进当前快照的清单里。生成一个新的Table Metadata JSON其中parent指向老版本同时把最新快照ID指向刚生成的新快照。把Catalog中该表的元数据指针从“老版本”原子地切到“新版本”。到这里Iceberg的一致性核心就出来了真正“提交”这个动作就是第5步的指针切换。指针切换之前新数据文件对任何读任务都不可见指针切换之后老版本对应的文件也不会立刻被删仍然能被老快照的读任务使用。这带来一个很舒服的效果读写完全不互斥。写任务不用锁表读任务也不用担心读到半截数据因为每个读任务都绑定一个固定的快照版本。这就是Iceberg能支撑高并发的一个底层前提——大多数人把它叫MVCC多版本并发控制本质上Iceberg是拿“不可变文件元数据指针切换”组合实现了一套轻量的MVCC。2. 提交竞争的唯一裁决点Catalog的CAS与乐观锁2.1 为什么说Iceberg的锁粒度“轻”到极致如果你想在传统关系型数据库里保证高并发写入的一致性常规做法就是加行锁、间隙锁、表锁锁粒度越小冲突越小但管理锁本身的成本很高死锁检测和锁等待都是大麻烦。Hive的ACID方案里也用锁但它在高并发场景下经常出现锁等待超时、死锁需要人工清理的情况。Iceberg故意绕开了这条路线。它不需要在数据行上加锁也不需要在整个表上持锁它只在“元数据指针”这个单点上做一次Compare-And-SwapCAS。谁抢到了最新版本指针的更新权谁就提交成功谁发现自己的基线版本已经落后谁就进入冲突处理流程。这个设计本质上是乐观并发大多数情况下不阻塞只有提交那一瞬间存在竞争竞争失败再决定怎么处理。对于数据湖这种“写多读多、单个提交动作本身很轻”的场景乐观并发比悲观锁的吞吐量要高得多。2.2 不同Catalog的CAS实现差异同样是CAS不同Catalog后端的具体实现方式不太一样。这块在生产环境里的影响比很多人想象的要大因为不同实现意味着不同的性能特征和潜在坑点。Catalog类型CAS实现的机制需要注意的问题Hive MetastoreHMS通过HMS中的表参数如metadata_location字段做版本校验更新时用条件更新HMS本身可能成为瓶颈高并发提交时对HMS压力较大Hadoop CatalogHDFS通过检查version-hint.text文件和metadata文件的时间戳/内容做比较依赖文件系统的一致性和时钟极端情况下可能出现校验误判JDBC CatalogSQL条件更新UPDATE ... WHERE version ?依赖数据库行锁但只在提交瞬间短暂持锁REST Catalog服务端持有元数据指针通过条件请求/版本号实现CAS最适合跨语言、跨引擎场景但要求REST服务高可用生产环境里最常见的是HMS和Hadoop Catalog。Hadoop Catalog很多人觉得简单实际上它对时钟同步和文件可见性非常敏感如果所在集群的NTP有问题有可能出现校验时读取到的元数据文件还是旧的导致提交判断出错。这点建议在依赖HDFS做提交仲裁的团队提前排查。2.3 乐观并发冲突的经典复现路径CAS冲突并不玄学它的产生路径非常固定。我自己用一个简单的双写实验复现过一次这个过程任务A和任务B同时读到当前表的metadata版本指向V10。A、B各自向HDFS写入自己的数据文件生成了各自的Manifest和新的metadata都基于V10。A先执行提交把指针从V10切到V11成功。B执行提交时发现当前指针已经不是V10了而是V11于是触发冲突。在Iceberg的日志里对应异常一般是CommitFailedException: Cannot commit because the current table metadata is [uuid/version] while the new metadata is based on [uuid/version]遇到这个异常不用慌它恰恰说明一致性机制在正常工作。如果两边都提交成功那才说明出问题了。真正需要做的是决定冲突之后干什么是重试还是放弃。3. 冲突之后的走向retry、overwrite还是error3.1 三种提交冲突解决策略的适用场景Iceberg的表属性里有一组参数专门控制提交冲突处理分别是commit.retry.num-retries、commit.retry.min-wait-ms、commit.retry.max-wait-ms、commit.retry.total-timeout-ms以及提交策略相关的commit.retry.override之类的行为。在实际行为上冲突后有几种典型走向Retry重试默认行为。Iceberg会在冲突发生时重新读取最新的表metadata基于新版本重新生成新的提交。这个过程对appends这类“无中生有”的写入很友好因为新加的多个数据文件之间天然是可累加的重新基于新版本再生成一次Manifest就能解决问题。Overwrite覆盖相当于关闭CAS直接把元数据指针覆盖成自己生成的版本。这种模式只适合某些特殊情况比如你明确知道表上最近没有其他并发写入或者这是一张重建性质的临时表。生产环境的高并发表我不建议用overwrite做常规提交策略它会把别人已经提交的数据“覆盖没”。Error报错直接把异常抛给上层。适合对数据精确度要求极高的场景宁可任务失败重跑也不愿在冲突时自动重放。很多人在配置里只看到了retry以为retry是万能药。其实要把问题拆开看待——retry重放的是什么操作类型很关键。3.2 一次线上冲突排查的完整链路复现前文提到的凌晨事故完整链路其实非常有代表性。两个Flink实时任务一个从Kafka读订单数据一个从Kafka读物流数据同时写入同一张Iceberg表都开启了checkpoint每两分钟做一次commit。从表面看两分钟一次提交冲突概率应该不高但事实是由于两个任务的checkpoint对齐时间比较接近又在同一批并行度下执行Commit操作经常在同一秒内发生。排查过程是这样走的第一步先看异常堆栈确认是CommitFailedException且堆栈里有“current table metadata”相关的版本校验信息。这一步就把问题定位在“提交冲突”而不是“写入失败”。第二步看两边的并发度。其中一个任务并行度是20每个subtask负责一个Kafka分区的数据到commit时是taskmanager统一提交相当于一批数据一次性commit。另一个任务并行度是16提交间隔几乎一样。两边在凌晨这个低峰期数据量都不大提交时间窗高度重合。第三步调整参数ALTER TABLE db.logistics_orders SET TBLPROPERTIES ( commit.retry.num-retries10, commit.retry.min-wait-ms200, commit.retry.max-wait-ms3000, commit.retry.total-timeout-ms1800000 );同时在上游任务里给提交动作加了一个30到60秒的随机抖动让两个任务的提交时刻错开。参数改动后冲突率从肉眼可见的频繁直接降为0。这个案例说明很多“高并发一致性”问题真正要调节的其实不是一致性机制本身而是提交节奏的合理性。高并发不是问题高并发且高频提交才是问题。3.3 Overwrite策略为什么被我长期禁用可能有人会问Iceberg官方保留了覆盖提交策略为什么这么排斥它我解释一下overwrite的实际效果当检测到当前metadata版本比事务基线新时overwrite策略不会重新读取最新版本做合并而是直接把你手里的旧版本以及对应文件集设置成表的最新状态。假设A刚提交了一批1000条的数据B手里是冲突前的metadataB用overwrite提交成功后A刚提交的1000条数据在表里“消失”了——因为B的metadata里不包含A的文件引用。A的数据文件还在存储里躺着但已经变成了孤儿文件除非手动清理否则就是一份白白占存储空间的脏数据。这太危险了。所以在团队内部我对生产表的要求是宁可让任务因为冲突失败后自动重跑也不要开overwrite。如果某些临时报表场景实在需要overwrite也必须隔离在独立的临时库中。4. 读侧为什么永远不慌快照隔离带来的“时间旅行”福利4.1 Snapshot是什么一个表的“CtrlZ存档”Iceberg的每次提交都会产生一个快照Snapshot。快照的本质是这个时间点上整张表所有数据文件的完整描述。每次快照在Table Metadata里按顺序排列记录了Snapshot ID、时间戳、操作类型append/overwrite/delete等以及对应的Manifest List位置。这个机制带来的直接好处就是“时间旅行”。你可以随时基于任意历史快照做查询前提是快照还没过期、底层数据文件还没被清理。晚高峰写入再频繁读任务只要绑定了一个快照ID就不会被后续的任何Commit影响。在Spark里查询历史快照很简单-- 基于时间戳查询 SELECT * FROM db.logistics_orders TIMESTAMP AS OF 2024-12-15 10:00:00; -- 基于快照ID查询 SELECT * FROM db.logistics_orders VERSION AS OF 8739812093812038;在Flink流读场景里可以通过starting-strategy参数控制要读取的快照位置比如earliest从最老可用快照开始latest只读新提交的数据。4.2 并发读时会发生什么先记“快照ID”再读文件读任务在执行时会先从Table Metadata中拿到当前已选定的快照ID然后顺着这个ID找到对应的Manifest List再找到所有相关数据文件。整个过程不会去管表的最新状态是什么。举个例子读任务R在10:00:00启动此时表最新快照是S5R绑定了S5。10:00:01时任务W提交了S610:00:02时任务W2又提交了S7。对R来说它从头到尾只认S5的文件清单S6和S7完全不影响它。这看起来没什么技术含量但在实际生产里价值极大。批流一体之所以能成立本质上就是靠快照隔离——批任务跑T-1的数据流任务读增量读到的是不同的快照区间彼此完全隔离不存在“读到别人写了一半的数据”这种状态。4.3 ExpireSnapshots和并发写冲突的隐藏坑快照隔离虽然好但它也有代价——快照会占存储空间。每个快照都保留了底层文件引用虽然元数据不大但底层文件因为多个快照引用而不能被删除存储占用会逐渐膨胀。因此常规维护任务是定期执行expire_snapshots把过期快照清理掉。这个维护动作有一个并发隐患如果清理快照的同时有个长时读任务正在用某个老快照扫描数据而这个快照恰好被expire了读任务就会因为底层文件被删而报FileNotFound异常。虽然Iceberg在Scan开始时会持有文件引用但在某些引擎的缓存策略下依然可能出现边界问题。我的实践建议是expire_snapshots统一在业务低峰期执行。保留至少72小时3天以内的快照spark.sql.catalog.mycat.warehouse下运行清理命令时设置older_than时间参数要留足余量。如果业务上有按小时回溯的需求快照保留窗口要按最极端回溯需求来设置而不是按默认值。CALL sys.expire_snapshots( table db.logistics_orders, older_than TIMESTAMP 2024-12-12 00:00:00, retain_last 5 );这个SQL的意思很直接清理2024-12-12之前的所有快照但至少保留最近5个快照。5. 生产环境踩过的五个一致性坑5.1 多个Catalog客户端缓存不一致导致“看不到刚提交的数据”现象很诡异任务A提交成功任务B读同一张表却一直看到旧数据持续了好几分钟才恢复。这不是Iceberg一致性出了问题而是Catalog客户端缓存导致的。在Spark里配置Iceberg Catalog时如果开了元数据缓存比如spark.sql.catalog.mycat.cache-enabledtrue spark.sql.catalog.mycat.cache-expiration-ms300000那么在缓存有效期内Spark可能不会重新拉取最新的metadata_location。对于实时性要求高的场景这个缓存很坑。我的做法是对实时入仓的作业统一关闭Catalog缓存或者把过期时间压到10秒以内。批任务可以容忍5分钟缓存流任务不能忍。5.2 并发compaction和正常写入互相踩Compaction小文件合并是数据湖的日常操作它本质上也是“读老快照 生成新数据文件 提交一个新快照”。如果正好赶上业务高峰期compaction任务和正常写入任务同时提交非常容易触发CAS冲突。冲突后如果compaction任务用的是retry它会基于最新快照重新处理可能把刚提交的文件又纳入合并范围导致白算一遍甚至产生重复合并的中间文件如果compaction任务反复失败最坏情况是合并任务一直无法推进小文件越来越多形成恶性循环。我的建议是两条腿走路业务高峰期不跑compaction专门在凌晨低峰期单独跑。compaction任务保证单实例运行不要在同一个表上同时起多个compaction作业。5.3 RemoveOrphanFiles误删正在使用的文件孤儿文件清理是另一个维护利器它会把“没有被任何元数据引用的文件”删除。正常情况下这是安全的因为提交有先后顺序已经提交的文件一定会被元数据引用。但在极端的并发场景下任务A刚写入数据文件还没完成提交如果这时候任务B执行了孤儿文件清理且清理的时间窗口覆盖了任务A的写入路径这些“还没被元数据引用”的文件可能被当成孤儿删掉导致任务A提交时找不到文件而失败。虽然这个窗口期很短但既然遇到了就不能忽视。操作上要满足两个条件再清理older_than参数必须早于当前正在运行的所有写任务的最早可能写入时间给足够缓冲。确认当前无正在提交中的任务。5.4 高并发小文件场景下元数据膨胀高并发写入经常带来另一个隐性成本——元数据膨胀。每次commit都会生成一个新的metadata JSON和Manifest List单次提交的数据量越大元数据文件增长越慢但如果是高频小批量写入一天几万次commit元数据目录的文件数量和总体积会以出乎意料的速度膨胀。解决思路有这几层控制提交频率尽量让单次提交的数据量接近合理大小比如在Flink里让checkpoint interval和Iceberg提交频率配合好。通过write.target-file-size-bytes控制单个数据文件大小避免产生过多小文件。定期做RewriteDataFiles和RewriteManifests来整理文件布局、合并小文件。5.5 随意提升retry次数带来的“雪崩效应”遇到提交冲突就调大重试次数这是不少人的第一反应。但在高并发高负载场景下无脑提高retry次数可能把问题放大。原因是这样的每个冲突任务在retry时都要先扫描最新metadata重新规划文件清单然后重新生成Manifest。如果整个集群同时有一大批任务在重试对HMS或文件系统造成的压力会反过来加剧故障——Catalog响应变慢提交超时触发更多重试更慢然后雪崩。我个人的经验数值是常规表commit.retry.num-retries设置为4到6次足够高并发实时入仓表最多调到10次超过10次基本说明写入链路本身有病要么是提交频率过高要么是应该拆分表。同时一定设置commit.retry.total-timeout-ms把整个重试过程的总时间卡死避免一个任务卡在提交环节几小时不退出。尾声一致性机制的边界不在机制本身没有万能的一致性方案。Iceberg的高并发一致性本质上是“乐观并发元数据原子切换快照隔离”这套组合拳的产物它很适合数据湖这种写入模式以append为主、读多写多、且能容忍一定冲突重试的场景。但如果你的业务需要的是行级强一致、高频随机更新Iceberg的Copy-on-Write delete/update会让写放大非常严重Merge-on-Read场景也需要更谨慎地设计compaction策略。我在线上跑了一年多之后最大的体会是机制只能保证“提交那一刻”的一致性而真正的端到端一致性要靠调度节奏、维护策略和参数配置共同配合。把提交节奏调平滑把冲突重试控制在合理范围把快照生命周期管好这张表才能在高并发下既稳又准。
阅读完成 · 觉得有帮助?
咨询建站