Apache Paimon 为什么能在流式数据湖场景中快速崛起

从架构与工程视角解析 Apache Paimon 为何能在流式数据湖场景快速崛起,介绍 LSM-Tree 存储、Changelog 增量输出与合并引擎等核心机制,对比 Iceberg、Hudi 的差异,并给出适用场景与落地建议。

实时链路里最别扭的那一段

过去几年,实时数仓的架构演进出现了一个很有意思的交叉点:流计算本身早就不是瓶颈了,但数据湖这一侧,一直缺一个能稳定承接流式写入、又能让下游按流或按批自由读取的存储层。

AI technology illustration

不少团队实际的链路长这样:Flink 实时清洗后的结果先落 Kafka,再由另一个作业抄写到 Hive 或 Iceberg。Kafka 负责短期存储和增量分发,Hive/Iceberg 负责长期存储和批量分析。这套链路能跑,但很绕——多一套消息队列就多一份延迟、多一份运维成本,还要处理两套存储之间的数据对账。

Apache Paimon 就是在这个节点上快速冒头的。它不是一个新计算引擎,也不打算替代数据库,而是一个把流式更新能力做进表格式(Table Format)里的流式数据湖存储方案。它的前身是 Flink Table Store,2023 年进入 Apache 孵化器后更名为 Apache Paimon,社区热度涨得很快。很多原来在 Flink 与 Iceberg/Hudi 之间做组合的团队,开始把 Paimon 放进候选名单,甚至愿意为它重新规划实时数仓的存储架构。

这篇文章不准备罗列特性清单,而是想从工程视角拆一下:Paimon 到底解决了链路里的哪些真实痛点,它的核心设计为什么很适合流式数据湖,以及什么样的团队、在什么阶段应该认真考虑引入它。

先把流式写入变成数据湖的一等公民

在聊 Paimon 之前,值得先理解一个问题:为什么数据湖做流式写入这么难?

传统 Hive 表是批的产物,文件按分区组织,写入要么重写分区,要么不断追加小文件。流式作业每隔一个 checkpoint 就落一批数据,直接写 Hive 的话,一小时后目录里可能躺着上千个小文件,查询性能急剧下降,还得额外安排合并任务来整理。

Iceberg 和 Hudi 分别从不同角度缓解了这个问题:Iceberg 用 manifest 管理文件的元数据,写路径仍然偏批;Hudi 的 MOR(读时合并)支持主键更新,但设计初衷更接近数据库的 upsert 语义,真正的流式语义是后加的能力。

Paimon 的选择更彻底:在存储层直接用 LSM-Tree(Log-Structured Merge Tree)来组织数据。写入先落在内存的 MemTable 中,达到阈值后刷成不可变的排序文件,后台的 compaction 再把小的排序文件合并成更大的文件,同时清理旧版本。这个设计在 KV 存储和时序数据库里很常见,天然适合高吞吐写入和频繁更新。

对主键表来说,LSM 带来的直接好处是:写入不再需要定位旧文件并覆盖,而是不断追加新的排序文件,读写之间不做原地修改,既避免了随机 IO,也让写入吞吐基本不随数据量增长而明显下降。真正麻烦的地方在读——同一主键的数据可能散落在多个排序文件里,Paimon 的读取端要做一个合并读(merge reader)来还原最新版本。

这也是很多刚接触 Paimon 的人容易困惑的点:为什么 Paimon 流式读带上更新之后,消费延迟比普通消息队列高?因为读端需要跨文件合并,还要兼顾批读的快照语义。理解了这层取舍,就不会拿它和 Kafka 做无意义的比较。

Changelog 直接落盘:增量不再是附加项

如果只是写入快,Hudi 其实也能做到一部分。Paimon 真正让流式场景兴奋的,是它能为主键表产出 changelog,也就是行级变更日志,包含新增(+I)、更新前(-U)、更新后(+U)和删除(-D)四种行类型。

有了 changelog,下游 Flink 作业就能像消费 Kafka topic 一样从数据湖里做真正的流式读取,而不是靠周期性的批扫描去假装有增量。这个能力在实时数仓里意义很大:以前增量数据要等实时链路算完再落湖,现在湖本身就能持续产出增量。

changelog 的产出方式可以配置,常用的有两种。一种是 input,在写入阶段直接记录原始输入带来的变化,成本低,但只反映这次写入了什么,不反映合并后的最终状态;另一种是 full-compaction,在触发全量合并后基于最终结果生成 changelog,输出稳定、完整的变更流,代价是合并成本更高,下游延迟取决于合并频率。

建一张带流式增量的主键表,DDL 大概长这样:

CREATE TABLE dwd_order_info (
  order_id      BIGINT PRIMARY KEY NOT ENFORCED,
  user_id       BIGINT,
  order_status  STRING,
  pay_amount    DECIMAL(10, 2),
  update_time   TIMESTAMP_LTZ(3)
) WITH (
  'bucket' = '4',
  'changelog-producer' = 'full-compaction',
  'full-compaction.delta-commits' = '20',
  'merge-engine' = 'deduplicate'
);

这里的 full-compaction.delta-commits 表示每累计多少次提交触发一次全量合并。下游逻辑依赖比较实时的数据,可以把间隔调小;更在意湖上查询性能,就把间隔调大。没有绝对正确的值,只有符合业务时效的取舍。

合并引擎下沉:多流打宽表的另一种解法

实时数仓最常见的场景之一,是把分散在几条流里的信息拼成一张宽表。以前的做法是在 Flink 里做 interval join,再 sink 到 MySQL、ES 或 ClickHouse。这个方案的痛点在于状态管理:join 状态的 TTL、乱序等待策略、并行度调整,每一项都需要专门调优。

Paimon 提供了一种新思路:让存储层直接按主键合并多条流的写入,这就是合并引擎(merge engine)。目前常用的有四种:

合并引擎 行为 典型场景
deduplicate 同主键后写覆盖前写 ODS 明细去重、CDC 主键幂等写入
partial-update 每条流只更新自己的字段 多流关联拼宽表,减轻 Flink 状态压力
aggregation 对数值字段做 sum/max/min 等聚合 实时指标累加、去重计数
first-row 只保留第一个到达的值 取首条有效信息,不被后写覆盖

拿 partial-update 举例。假设要构建订单域的实时宽表,订单流、支付流、物流流来自三个独立的 Kafka topic,原来的做法是在 Flink 里维护一个非常大的 keyed state。用 Paimon 之后,三条流分别往同一张主键表写各自的字段,存储层按 order_id 合并,下游再通过流式读订阅合并后的结果。

-- 支付流写入
INSERT INTO dwd_order_broadcast (order_id, pay_amount, pay_time)
SELECT order_id, amount, pay_time FROM pay_topic;

-- 物流流写入
INSERT INTO dwd_order_broadcast (order_id, ship_status, ship_time)
SELECT order_id, status, ship_time FROM ship_topic;

这个方案最直接的价值,是把 Flink 的大状态换成了数据湖里的文件,状态管理从任务的内存和 RocksDB 转移到了存储层的合并逻辑,重启恢复和调整并行度都不再受制于 join 状态的大小。代价是合并时效取决于 compaction 策略,不是严格的毫秒级。

还要提醒一句:partial-update 不是万能的。如果两条流对同一字段都有写入,先后顺序会影响最终结果,需要结合字段的业务含义设计 sequence field 或明确时序约定。

和 Iceberg、Hudi 比,差异到底在哪

很多团队选型时会同时考虑 Paimon、Iceberg 和 Hudi。这三者严格来说不完全是同一层面的东西,但在湖仓场景里经常被放在一起权衡。我用一个简表来梳理它们各自的定位:

方案 核心设计 流式语义 典型定位
Iceberg 表格式 + 快照隔离 偏批,增量读取有限 离线湖仓的统一表格式
Hudi MOR + Upsert 支持增量,运维较重 数据库入湖、更新密集场景
Paimon LSM + Changelog 原生流式读写 实时湖仓、Flink 生态数据处理

更具体地说,Iceberg 的优势在于对存量大数据生态的兼容性和 ACID 保证,适合以批为主、辅以增量的湖仓底座。Hudi 的 MOR 模型确实能处理主键更新,但小文件管理、compaction 调度、流式读的一致性保障都需要额外的工程投入。Paimon 则从底层就把流作为第一设计目标,Flink 作业写 Paimon 的门槛明显更低,流批一体的体验也更顺。

不过这里要补一句公道话:Paimon 的快速崛起有相当一部分来自 Flink 生态的红利,它在 Spark、Trino 等引擎上的成熟度仍在追赶中。如果团队的计算引擎以 Spark 为主,Paimon 的吸引力会打不少折扣。

常见的几个误区

随着热度上升,社区里也出现了一些被过度简化的说法,挑三个最常见的:

  • Paimon 能替代 Kafka。不能。Paimon 的数据可见性通常有 checkpoint 级到分钟级延迟,适合做湖上的增量分发,不适合做毫秒级消息投递。消息队列和湖存储是上下游关系,不是替代关系。
  • Paimon 是 Flink 专属。不准确。Spark、Trino、StarRocks、Doris 等引擎都能接入 Paimon,只是 Flink 生态的集成最成熟,很多能力是 Flink 版先落地,其他引擎后跟进。
  • 外部表也能做更新。外部表对应 Hive 兼容的纯追加场景,不走 LSM 和主键合并,把它当主键表用会得到错误结果。用之前一定先分清主键表和外部表的边界。

这些误区背后其实是同一个问题:把 Paimon 当成无所不能的银弹。它只是把流式数据湖的问题解决得更优雅,并没有改变数据架构的基本约束。

什么时候值得引入 Paimon

如果团队正在经历下面这些信号,Paimon 值得认真评估:

  • 数据已经在 Kafka 和 Hive 两套体系里维护,对账成本越来越高;
  • 需要把整库 CDC 数据实时入湖,并且下游要订阅变更;
  • 实时宽表依赖 Flink 大状态 join,任务稳定性被状态拖累;
  • 希望同一份湖上数据既能跑批分析,又能被流式任务持续消费。

落地时建议从小场景切入,不要一上来就做全景湖仓。最自然的起点是替换 Kafka 抄写到 Hive 这段链路,先用 append 表承接实时落盘,让自动 compaction 处理小文件问题;跑通后再升级到主键表,逐步把 changelog 和合并引擎用起来。

另一条有效路径是 Flink CDC 整库入湖。用现成的整库同步能力把业务库变更实时写入 Paimon,再让下游 StarRocks、Doris 或 Flink 作业消费。这个场景里 Paimon 的收益非常直接,省掉了自研 CDC 管道的大量重复工作。

关于 bucket 数需要专门提醒:它决定数据的分桶粒度,直接影响写入并行度和读取并发。太小容易热点倾斜,太大容易产生文件碎片,而且后期调整成本不低。建议先按数据量和更新频率估算,再通过压测验证。

最后说两句

Paimon 能在流式数据湖场景里快速崛起,表面上看是社区运营和 Flink 生态的推动,本质上还是因为它踩中了一个长期存在、但一直没被很好满足的需求:数据湖需要真正懂流的存储底座。LSM-Tree 提供了高吞吐写入的基础,changelog 让湖具备增量输出的能力,合并引擎则把复杂的流式关联逻辑从计算层下放到了存储层。

当然,任何技术都有适用边界。Paimon 的 compaction 成本、跨引擎成熟度,以及它在超大规模生产环境下的长期表现,还需要更多案例来验证。但至少在今天,它已经成为实时湖仓选型里绕不开的选项。如果你的团队正在为流和批到底能不能统一而纠结,花两周时间做一个 PoC,用真实数据和真实链路去验证,比看任何对比文章都更有说服力。

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

(0)
上一篇 41分钟前
下一篇 34分钟前

相关推荐