湖仓一体架构中如何统一批处理和流处理的读写语义

本文解析湖仓一体架构下批处理与流处理读写语义的统一方式,从事务日志、快照读、增量读和幂等写入等关键机制展开,对比Delta Lake、Iceberg、Hudi的实现差异,并给出工程落地建议。

湖仓一体架构的讨论已经持续了好几年,很多团队也确实把它当成新一代数据基础设施的方向。不过真正落地之后,不少人会碰上一个很具体的问题:批处理和流处理明明操作同一份存储,读出来的结果却经常“对不上”。离线任务和实时任务都跑在同一张表上,但各自的边界、可见性和重复处理能力完全不同,于是问题就会从“数据怎么迟迟没更新”变成“数据读出来怎么跟我预期不一样”。

AI technology illustration

要解决这个问题,得先从读写语义入手。批处理读数据通常是对某个固定数据集合做一次全量扫描,而流处理面对的是无界数据,天然需要持续读取、增量处理。湖仓一体架构通过事务日志和版本快照,试图把这些差异收敛到底层存储中。这篇文章会从读和写两个角度,拆解统一语义背后的关键机制,也聊一聊不同表格式的取舍。

批处理和流处理在读写语义上到底差在哪

批处理和流处理的差异,不仅仅是实时和离线的时效差别。它们对“数据边界”的理解就不一样。批处理任务启动时,就假定数据集合已经固定,读完之后哪怕源表再发生变化,本次任务的结果也不会变。流处理则没有这个假设,数据源源不断地进来,任务必须自己记住“已经读到哪了”,才能在下一次启动时续上。

这个差异一旦落到同一张湖仓表上,就开始变得拧巴。一个典型的场景是:Flink任务每十秒往表中提交一批数据,Spark批处理在任意时刻启动并尝试读取全表。如果存储层没有统一的事务边界,Spark可能读到本批次写入了一半的文件,产生不完整视图。即使底层存储能保证读快照,频繁的提交也会让批处理在解析元数据时付出额外开销。

所以这里的关键不是“快”和“慢”,而是数据从写入到可见是否具有清晰、稳定的事务语义。批处理希望每次都看到一个一致的数据版本,流处理则希望每次写入都不阻塞、不带来太大元数据压力。两者的诉求,最终都指向同一个点:我们需要一个所有读写都能共同遵守的提交协议。

事务日志:统一读写语义的桥梁

在湖仓一体出现之前,数据湖上的写入通常就是直接生成、删除文件。没有提交记录,也就没有版本概念。流处理可以通过Kafka offset记住自己的进度,批处理却不知道哪些文件是新的,哪些文件是旧版本留下的。数据表变得不可读,也几乎无法做增量消费。

湖仓一体的表格式,比如Delta Lake、Iceberg、Hudi,都引入了事务日志(或等价物)。每次写入,不管是批处理一次性覆盖,还是流处理小批量提交,都会在事务日志中留下一个原子提交记录。读取端通过日志来获取表的当前版本和所有历史版本,而不是直接去目录里列文件。

还是用上面的场景:Flink每次提交一批数据,都会追加一条提交记录,Spark批处理启动时读到最新提交对应的快照,就能看到一个完整的数据视图。即使两个任务并行执行,也没有读取半批次的风险。更重要的是,这条提交日志同样可以被流处理利用。流处理不必从头扫表,只需要跟踪新出现的提交记录,就能拿到增量数据。这样一来,读端语义就统一了:批处理读完整快照,流处理读增量,但两者都遵循同一套提交协议。

统一读语义:从快照读扩展到增量读

有了事务日志,批处理读快照这件事变得很自然。真正的挑战在于如何让流处理也有类似语义。

以Delta Lake为例,Spark Structured Streaming可以直接从一个指定的版本或时间戳开始,读取后续提交中出现的变更:

spark.readStream \
  .format('delta') \
  .option('startingVersion', '120') \
  .table('events')

如果用时间戳,则用startingTimestamp参数。这个读取过程会分析事务日志中版本120之后的每次提交,把新增文件作为流的新一批输入。Iceberg的Flink连接器也有类似能力,基于snapshot-id或时间戳扫描增量快照。Hudi的增量查询同样是围绕Timeline上的commit记录展开。

不过,这种增量读并不是免费的能力。它有一个隐藏成本:流处理需要经常扫描事务日志和元数据。如果写入端提交频率太高,比如每5秒一次,增量读就会被大量小版本淹没,有时甚至比直接读文件还要慢。实践中建议将流写入的commit间隔控制在10秒到1分钟之间,并根据延迟要求重新设计微批的触发时间。把它当成一个整体来调优,而不是只调读端。

统一写语义:幂等和Merge是关键

读语义解决的是“看到什么”,写语义解决的是“怎么写入才能不被流处理的重试弄脏”。流处理因为在分布式环境下运行,失败重跑几乎是必然的。如果写操作不幂等,那么Kafka消息重放几遍,表里的数据就会重复几遍。

这时可以把批处理中常用的merge操作搬到流处理中。其示意如下:

MERGE INTO user_profile AS t
USING (SELECT * FROM streaming_updates) AS s
ON t.user_id = s.user_id
WHEN MATCHED THEN UPDATE SET t.last_active = s.last_active
WHEN NOT MATCHED THEN INSERT *

这段SQL假设流处理已经将更新数据写到临时视图streaming_updates中。无论流任务重放多少次同样的消息,每个user_id最终都会只保留最后一次更新的结果。这就是幂等写入。Delta Lake、Hudi都支持类似操作,Iceberg目前对merge的支持在不同引擎里也有所覆盖。

这里有一个很重要的经验:一个流任务如果只是append日志型数据,相对安全;一旦要更新已有行,就必须定义主键和更新策略。否则,写语义不统一,批处理读出来的分析结果在重复数据下很难让人信任。

另一个写语义的坑是小文件。流处理每个微批都会生成若干文件,如果提交间隔很短,表里的小文件数量会快速增长。此时下游批处理的元数据扫描、compaction和查询都可能被拖慢。通常做法是定期执行压缩,例如Delta的OPTIMIZE,Hudi的clustering,Iceberg的rewrite data files。同时,控制写端的提交频率,比事后无限压缩更有意义。

三种主流表格式的读写语义支持对比

表格式 事务日志 增量读语义 写语义常用方式 典型引擎配合
Delta Lake DeltaLog 从指定版本/时间戳流读 append、merge、CDC Spark Structured Streaming、Flink
Apache Iceberg Snapshot Metadata 基于snapshot-id/时间戳增量扫描 append、overwrite、merge Flink、Spark
Apache Hudi Timeline 基于commit/increment增量读取 upsert、insert、bulk_insert Spark、Flink

这张表不是让你做选型决策,而是为了看清不同表格式对读写语义的侧重点。Delta Lake和Spark生态绑定比较深,增量流读在Structured Streaming里体验完整。Iceberg的元数据设计更通用,多引擎支持也做得好,如果你本来就在用Flink做实时计算,可以考虑它。Hudi则把upsert当作核心能力,适合存在大量主键更新的场景,比如用户画像表。

注意,表格式的成熟度不是静态的,选择时更要看团队已有的计算引擎、监控体系和数据建模习惯。

三个容易踩的误区

关于统一批流读写语义,在实践中见过不少团队走弯路,以下几个误区最常见。

  1. 把“批流一体”理解成“同一套作业同时跑两种模式”。其实批处理和流处理的计算模型仍然有很大差别,不可能靠一个框架完全抹平。湖仓一体统一的是数据表层的读写语义,不是计算逻辑。
  2. 忽视幂等写入。真正跑生产的时候,一个流任务因为网络抖动或资源屏蔽导致失败,重启后从checkpoint恢复,但Kafka的位点和表状态之间没有原子关联。这时候假如写链路是简单append,重复数据就会进入表。
  3. 没有统一的时间语义。批处理可能用“事件发生日期”这样一个字段来圈定分区,流处理却用Flink的watermark来决定窗口触发。如果表里没有严格约定标准事件时间,批和流对同一批数据的判断就会错位。

落地时,可以从这三件事开始

如果你所在团队正准备在湖仓一体中统一批流读写,不用一开始就推倒重来。建议先从以下三件事做起。

  • 统一核心业务表的表格式。所有需要批流共同访问的表,至少使用同一种支持事务日志的表格式,避免一部分表是普通Hive表、另一部分是Delta表,导致批流语义不一致。
  • 为流式读取定义规范。所有实时任务需要读取湖仓表的新增数据时,明确使用所在表格式的增量读API,并统一指定起始版本或时间戳。
  • 建立写入元数据的监控。监控项包括小文件数量、提交频率、compaction拖载情况、merge冲突率。不要等到出现百万小文件或数据重复后再处理。

另外一个实际场景是,有些团队早期为了快速上线实时链路,直接让Flink写HDFS普通文件,然后Spark批处理定时去目录里发现新文件。也许两周内能跑通,但一旦需要处理更新、回放、增量,这套方案就彻底卡住。后来迁移到Delta表,用了增量读和merge,反而把架构简化了。这个例子想说明的是,读写语义的统一,本质上是为了让你不需要设计两套完全不同的数据访问协议。

统一读写语义不是设置一个参数就能完成,而是一套从数据模型、存储格式到任务容错的完整约定。

最后再回到开头的问题。湖仓一体架构中统一批处理和流处理的读写语义,并不是为了强行让批处理和流处理变得一模一样,而是让它们能够在同一份数据上健康共存。批处理依然可以用它喜欢的方式做全量分析,流处理依然可以持续低延迟地写入,关键是它们都通过同一个事务日志、同一套数据版本来沟通,而不是各自为政。

理解这一点,你可能就不会再纠结于“为什么批流一体还是有两个API”这类问题。读写语义统一,指的是存储层给出的承诺统一,而不是API统一。它能让你在需要批处理能力的时候放心做全量计算,在需要流处理能力的时候安全地增量更新。如果这个过程里遇到任何怪异的数据问题,不妨先回头查一查:提交协议是不是一致的、写入是否幂等、时间语义是否统一。大概率答案就藏在这三件事里。

原创文章,作者:fudengji,如若转载,请注明出处:https://fudengji.cn/article/978/

(0)
上一篇 16分钟前
下一篇 27秒前

相关推荐