1. 从“数据仓库”到“数据湖”:为什么我们需要Paimon?

如果你在过去几年里处理过大数据,大概率经历过这样的场景:业务部门临时需要一个报表,你发现数据源在MySQL,一部分历史数据在Hive,实时数据又跑在Kafka里。于是你开始写Spark作业,吭哧吭哧地做ETL,把数据统一格式、清洗、合并,最后导入到一个新的Hive表里。几天后,报表需求变了,你发现之前合并的逻辑有问题,或者需要回溯某一天的历史数据,但原始Kafka数据因为保留策略已经被清理了。这种“数据孤岛”和“数据回溯”的痛,是传统数仓架构难以根治的顽疾。

数据湖的概念就是为了解决这些问题而生的。它不像数据仓库那样,要求数据在入库前就必须有严格的结构(Schema-on-Write),而是允许你以原始格式(如Parquet、ORC、JSON)先把海量数据“倾倒”进来,在查询时再定义结构(Schema-on-Read)。这带来了极大的灵活性,但也引入了新的挑战:如何管理这些海量文件?如何保证流批数据的一致性?如何高效地更新和删除记录?如何支持时间旅行查询(Time Travel)?

这正是Apache Paimon(原名Flink Table Store)要回答的问题。它不是一个简单的文件存储格式,而是一个 构建在数据湖存储之上、具备数据库表管理能力的存储层 。你可以把它理解为一个“湖仓一体”(Lakehouse)的实现。它底层使用对象存储(如S3、OSS)或HDFS来存放文件,但在上层提供了像数据库一样的 ACID事务、主键更新、增量读取和流批统一访问 的能力。这意味着,你可以用流的方式持续写入数据,同时用批的方式做历史分析,并且两者看到的数据状态是完全一致的。

简单来说,Paimon试图让你用管理数据库表一样简单的方式,去管理存储在廉价对象存储里的海量数据,同时兼顾流处理和批处理的效率。接下来,我们就剥开它的外层,看看它是如何实现这一目标的。

2. Paimon架构总览:三层抽象与两种“表”

要理解Paimon,首先要摒弃“它只是一个文件格式”的想法。它的架构可以抽象为三层: 存储层、表管理层和计算层

存储层 是基石,就是你的对象存储(S3、OSS)或分布式文件系统(HDFS)。Paimon的所有数据文件(Data Files)、清单文件(Manifests)和快照(Snapshot)都物理存储在这里。它的选择决定了存储的成本和持久性。

表管理层 是Paimon的核心大脑。它负责将上层(计算层)的“插入”、“更新”、“删除”等操作,翻译成对底层文件系统的“新增文件”、“标记删除”等动作,并维护一套元数据来保证数据的一致性视图。这层的关键组件是:

  • Snapshot(快照) :这是Paimon实现ACID和时间旅行的核心。每次提交(比如一批流数据写入完成)都会生成一个新的快照。快照是一个指向当前所有有效数据文件的指针列表。查询时,只需读取最新快照对应的文件,就能获得一致的数据视图。回溯历史?只需指定一个历史快照ID即可。
  • Manifest(清单) :快照不能直接指向成千上万个数据文件,那样效率太低。清单文件就是快照和数据文件之间的索引。一个快照会引用一个或多个清单文件,每个清单文件则记录了属于该快照的一批数据文件的路径、统计信息(如最小值、最大值)等。
  • Data File(数据文件) :实际存储用户数据的文件,默认采用列式存储格式ORC或Parquet,以提供高效的分析查询性能。

计算层 是Paimon的“手足”,负责数据的读写。它完美集成了Apache Flink,作为其原生的一等公民,可以通过Flink SQL、DataStream API或Table API进行流式读写和批处理。同时,它也支持通过Spark、Hive、Trino/Presto等引擎进行批查询,实现了计算引擎的解耦。

在表管理层,Paimon设计了两种核心的“表”类型,对应不同的数据变更处理模式,这是理解其原理的关键:

2.1 Changelog表:流式更新的核心

这是Paimon的默认表类型,也是其流批一体能力的体现。这种表要求你定义主键(Primary Key)。它的核心思想是: 将所有的数据变更(插入、更新、删除)都转化为带有“增/删”标记的记录(Changelog),并有序地存储下来。

想象一下数据库的Binlog,Paimon的Changelog表就在做类似的事情。当你执行一条 UPDATE 语句时,Paimon不会去物理地修改已有的数据文件,而是会生成两条记录:一条 -D (删除)记录表示旧值的“删除”,一条 +I (插入)记录表示新值的“插入”。这些包含变更语义的记录会和其他插入的记录一起,被写入新的数据文件中。

为什么这么做?

  1. 简化流处理 :流计算引擎(如Flink)可以直接消费这些完整的Changelog流,无需复杂的“拉链表”或“全量+增量”合并逻辑,就能构建实时物化视图或更新下游维度表。
  2. 高效合并 :虽然存储的是变更日志,但Paimon在后台会运行一个名为 Compaction 的进程。这个进程会将多个包含大量更新日志的文件,合并成少数几个包含最终数据状态的文件,从而优化查询性能。这个过程对用户是透明的。
  3. 支持事务 :一次提交内的所有变更,要么全部生效(生成一个新快照),要么全部失败(快照不变),保证了ACID中的原子性和一致性。

实操心得 :在定义Changelog表时,主键的选择至关重要。它不仅是数据更新的依据,也直接影响Compaction的效率。主键字段应选择更新频率适中、区分度高的业务字段,如 order_id 。避免使用像 update_time 这种持续变化的字段作为唯一主键,否则会导致每次更新都被视为一条新记录,产生大量无效的 -D/+I 对,加剧文件膨胀。

2.2 Append-Only表:高性能写入的权衡

并非所有数据都需要更新。例如,日志数据、交易流水、IoT传感器读数,这些数据一旦产生就不会改变,只会追加。对于这种场景,使用Changelog表反而会引入不必要的主键约束和更新开销。

Append-Only表就是为此而生。它不要求定义主键,所有写入操作都被视为纯粹的 插入(Insert) 。它的架构因此变得非常轻量:

  • 无主键管理 :省去了维护主键索引和解决更新冲突的开销。
  • 更简单的Compaction :后台合并任务只需将小文件合并成大文件,无需处理复杂的更新合并逻辑,速度更快。
  • 更高的写入吞吐 :由于逻辑简单,写入路径更短,通常能获得比Changelog表更高的写入性能。

如何选择?

  • 如果你的数据有明确的更新需求(如用户画像、商品库存、订单状态),必须使用 Changelog表
  • 如果你的数据是只增不改的(如操作日志、点击流、流水记录), Append-Only表 是更高效的选择。在Flink SQL中,可以通过 'primary-key' = '' 不设置主键来创建此类表。

3. 核心原理深度拆解:快照、Compaction与索引

理解了两种表类型,我们深入到Paimon如何实现这些能力的细节中。三个核心机制构成了它的基石:快照隔离、Compaction合并和索引加速。

3.1 快照隔离与时间旅行:数据一致性的基石

这是Paimon区别于简单文件存储最核心的特性。它借鉴了数据库和多版本并发控制(MVCC)的思想。

工作原理

  1. 写入 :当一批数据写入完成并提交时,Paimon会创建一个新的 快照(Snapshot) 。这个快照本身是一个很小的元数据文件,记录了本次提交生成的 新增数据文件列表 ,并指向其父快照(上一次提交的快照)。
  2. 快照链 :所有快照通过父子关系形成一个链表。最新的快照被称为“ 当前快照(Current Snapshot) ”,它代表了数据集的当前完整状态。
  3. 读取 :任何查询在开始时,都会“锚定”到一个特定的快照(默认是最新快照)。查询引擎根据该快照找到其引用的所有清单和数据文件,读取这些文件,就能获得一个 在某个时间点上绝对一致的数据视图
  4. 时间旅行 :要查询历史数据,只需在查询时指定一个历史快照ID或时间戳。Paimon会定位到那个时间点的快照,并读取其对应的文件集合。因为历史数据文件从未被物理删除或覆盖,所以这个查询是可行且高效的。

带来的好处

  • 读写分离 :一个长时间运行的批处理查询可以锚定在开始时的快照,不受后续流写入的影响,保证了结果的可重复性。
  • 回滚与审计 :可以轻松地将表回滚到之前的任何一个健康状态,或者审计历史上任意时刻的数据内容。
  • 增量读取 :流处理任务可以持续消费从一个快照到下一个快照之间 新增 的数据文件(即增量数据),这是实现流处理的基础。

3.2 Compaction:化“日志”为“状态”的魔法

对于Changelog表,如果一直存储原始的 -D/+I 变更记录,查询性能会急剧下降,因为要扫描大量无效的中间状态。Compaction(压缩合并)就是那个在后台默默将“流水账”整理成“总账本”的管家。

Compaction主要做两件事

  1. 小文件合并 :将多次写入产生的小数据文件(读放大)合并成更大的文件,减少查询时需要打开的文件句柄数,提升I/O效率。
  2. Changelog合并(关键) :这是针对Changelog表的专属优化。它会扫描多个文件, 将属于同一主键的多条变更记录(如 -D, +I, -D, +I )合并成最终的状态记录(一条最新的 +I 或直接删除) 。合并后的文件只包含数据的最新状态,查询时无需再遍历变更历史。

Paimon的Compaction策略是可配置的。你可以设置触发Compaction的文件大小阈值、时间间隔等。在流写入场景下,通常建议开启 全异步Compaction ,让一个独立的Flink作业专门负责合并,这样不会阻塞主写入流,保证写入延迟的稳定。

踩坑记录 :Compaction资源不足是线上常见问题。如果写入流量很大而Compaction速度跟不上,会导致小文件和未合并的Changelog堆积,查询速度变慢,甚至最终因文件数过多导致元数据过大而影响稳定性。我们的经验是,为Compaction作业单独分配足够的CPU和内存资源,并监控 number-of-files changelog-size 这类指标。对于更新非常频繁的热点数据,可以考虑根据业务分区,避免全表大合并。

3.3 索引:加速查询的利器

虽然Paimon依赖计算引擎(如Spark)的文件过滤能力,但它也内置了索引机制来进一步加速点查和范围查询。

  • 主键索引 :对于Changelog表,主键是天然的索引维度。在Compaction合并文件时,Paimon会按主键对数据进行排序和索引(在ORC/Parquet文件内部)。查询时,引擎可以利用这些索引信息快速定位数据块。
  • 二级索引 :Paimon支持在非主键字段上创建二级索引(如Bloom Filter)。例如,为常用的查询过滤字段 user_id 创建布隆过滤器索引,可以在读取文件时快速跳过绝对不包含目标值的文件,大幅减少I/O。
  • 分区与分桶 :这是最常用且高效的“索引”。通过 PARTITIONED BY BUCKET 关键字,可以将数据物理地组织到不同的目录和文件中。
    • 分区 :常用于时间维度(如 dt='2024-05-20' ),直接利用HDFS的目录结构进行剪枝,对于按时间范围查询的过滤效果极佳。
    • 分桶 :根据主键或某个字段的哈希值,将数据分散到固定数量的桶文件中。这能保证相同键值的数据落在同一个文件里,对于点查和更新操作,可以避免全表扫描,只需读取一个桶文件。

一个典型的表定义会结合使用这些技术:

CREATE TABLE orders (
    order_id BIGINT,
    user_id BIGINT,
    amount DECIMAL(10,2),
    status STRING,
    dt STRING,
    PRIMARY KEY (dt, order_id) NOT ENFORCED
) PARTITIONED BY (dt)
WITH (
    'bucket' = '10',
    'index.bloom-filter.columns' = 'user_id'
);

这个表按天分区,每天的数据又哈希成10个桶。查询 dt='2024-05-20' and order_id=123 时,Flink或Spark能快速定位到 dt=2024-05-20 分区下的某个桶文件,并利用主键排序和布隆过滤器快速找到数据。

4. 典型应用场景与实战配置指南

理解了原理,我们来看看Paimon在哪些场景下能大放异彩,以及如何针对性地进行配置。

4.1 场景一:CDC入湖与实时数仓

这是Paimon最经典的应用。使用Flink CDC直接捕获MySQL、PostgreSQL等数据库的变更日志,实时写入Paimon表。

架构价值

  • 实时同步 :替代了传统的离线T+1数据同步,实现秒级延迟。
  • 流批统一 :同一张Paimon表,流处理作业可以消费其增量日志做实时聚合,批处理作业可以读取全量做T+1报表,数据同源一致。
  • 历史回溯 :任何数据问题都可以通过时间旅行查询定位到变更点。

关键配置

-- Flink SQL创建CDC源表并写入Paimon
CREATE TABLE mysql_orders (
    id BIGINT,
    ... ,
    PRIMARY KEY (id) NOT ENFORCED
) WITH (
    'connector' = 'mysql-cdc',
    ...
);

CREATE TABLE paimon_orders (
    id BIGINT,
    ... ,
    PRIMARY KEY (id) NOT ENFORCED
) PARTITIONED BY (dt)
WITH (
    'bucket' = '5',
    'changelog-producer' = 'full-compaction', -- 关键:确保生成完整的Changelog供下游消费
    'compaction.interval' = '1h',
    'snapshot.time-retained' = '7d' -- 根据业务需要保留快照
);

-- 执行写入
INSERT INTO paimon_orders SELECT *, DATE_FORMAT(update_time, 'yyyy-MM-dd') as dt FROM mysql_orders;

注意 changelog-producer 配置至关重要。 'full-compaction' 模式能保证在每次Compaction后,为每个主键生成完整的 -U/+U 变更流,这对于下游流作业正确计算非常重要。如果下游只需要最终状态,可以选用 'lookup' 模式以降低开销。

4.2 场景二:流计算结果的物化存储

在实时数仓的DWD或DWS层,经常需要将流式聚合(如每分钟的GMV、UV)的结果持久化存储,供即席查询或下游批处理使用。

传统痛点 :将结果写入Kafka,再通过批作业周期性地导入Hive,链路长、时效性差、一致性难保证。 Paimon方案 :直接将Flink流聚合的结果写入Paimon表。

优势

  • 查询即所得 :聚合结果一旦写入Paimon,任何支持Paimon的查询引擎(如Trino)都可以立即查询到最新结果。
  • 自带分区 :可以按时间(如 hour )分区,方便按时间范围快速查询。
  • 更新自然 :如果聚合逻辑涉及对历史数据的修正(如去重结果更新),Paimon的主键更新能力可以轻松应对。

配置要点 :此类表通常是Append-Only或低更新频率的Changelog表。应合理设置分区和分桶,并调大 compaction.file-size 以减少小文件,因为写入模式通常是规律的批量追加。

4.3 场景三:替代Hive,作为统一的离线数仓表层

对于已有的Hive离线数仓,可以考虑将核心的ODS层或DWD层表迁移到Paimon。

操作路径

  1. 使用Spark或Flink批作业,将历史Hive数据批量导入(Bootstrap)到Paimon表。
  2. 将新的增量数据(如每日调度任务产生的数据)直接写入Paimon。
  3. 将下游的Hive SQL查询,逐步改为用Spark/Trino查询Paimon表。

收益与挑战

  • 收益 :获得了更新删除能力、时间旅行、流批统一入口。
  • 挑战 :历史数据迁移成本、下游作业改造、生态工具适配(如数据质量检查工具)。建议从单点业务开始试点,验证稳定性和性能收益后再推广。

4.4 性能调优核心参数

要让Paimon跑得又快又稳,以下几个核心参数需要根据数据规模和模式进行调整:

  • bucket :分桶数。这是影响并行度和文件大小的关键。建议设置为写入并行度的整数倍,且每个桶最终的文件大小在128MB ~ 1GB为宜。太小则文件过多,太大则不利于并行。
  • compaction.interval :Compaction触发间隔。流处理中通常设置为 '1h' '30min' 。对于写入压力大的场景,可以缩短间隔,但会增加后台资源消耗。
  • changelog-producer :变更日志生成模式。 'full-compaction' (保障完整流,延迟高) vs 'lookup' (延迟低,可能丢失中间状态)。根据下游消费者需求选择。
  • snapshot.num-retained.min / snapshot.time-retained :保留的快照数量/时间。用于控制时间旅行的深度和存储成本。需根据业务审计需求设置。
  • scan.parallelism :批查询时的并行度。在Spark或Flink批查询中,合理设置此参数能充分利用集群资源,加速全表扫描。

5. 选型对比与生态定位:Paimon vs. Iceberg vs. Hudi

谈到数据湖格式,Apache Iceberg和Apache Hudi是无法绕开的对比项。三者目标相似,但设计哲学和适用场景各有侧重。

特性维度 Apache Paimon Apache Iceberg Apache Hudi
核心设计理念 流批一体原生 ,为Flink深度优化,强调作为流处理Sink和Source的体验。 通用表格式 ,定义开放的、引擎无关的Table Format标准,追求极致的兼容性和可靠性。 快速Upsert ,最初为Spark生态的增量更新和删除场景设计,强调低延迟的数据摄取。
流处理集成 原生最佳 ,与Flink API无缝集成,Changelog概念原生支持,流读写体验最自然。 通过Flink Connector支持,功能完善,但流式更新语义需要依靠“Equlity Delete”等实现,稍显复杂。 通过Flink Connector支持,但核心流式模型(如COW/MOR)最初为Spark设计,在Flink生态集成深度稍逊。
批处理查询 支持良好,通过Spark、Hive、Trino Connector。 支持极佳 ,拥有最广泛的生态支持(Spark, Trino, Presto, Impala等),是批查询生态最成熟的格式。 支持良好,主要通过Spark和Hive。
更新删除效率 主键更新,通过Compaction合并Changelog,适合高频更新。 通过“Copy-on-Write”或“Merge-on-Read”模式支持,MoR模式适合写多读少。 Upsert效率高 ,特别是“Merge-on-Read”模式,为频繁的插入/更新操作优化。
事务与时间旅行 基于快照,支持ACID和时间旅行。 基于快照,支持ACID和时间旅行,实现非常严谨。 基于时间线(Timeline),支持增量读取和有限的时间旅行。
学习与使用成本 在Flink生态中简单直观,概念较少。 概念相对较多(Snapshot, Manifest, Partition Spec等),但文档和社区成熟。 概念较多(COW/MOR, Index类型等),配置选项复杂。
典型适用场景 Flink为中心的实时数仓 ,CDC实时入湖,流计算结果物化。 企业级离线/湖仓一体 ,需要多引擎(Spark, Trino)稳定访问,对开放性和可靠性要求高。 近实时数据湖 ,对数据摄取延迟要求极高,有大量Upsert需求的场景(如Spark流处理)。

如何选择?

  • 如果你的技术栈以Apache Flink为核心 ,构建实时数据管道和流批一体数仓,追求极致的流处理开发体验, Paimon是目前最自然、最顺畅的选择 。它就像是Flink的“亲儿子”,很多流处理中的痛点(如Exactly-Once Sink、流读Changelog)都被原生解决了。
  • 如果你需要一个稳定、开放、被众多查询引擎广泛支持的企业级表格式 ,用于整合公司内多样化的计算工具(Spark, Presto, Impala, Athena等), Apache Iceberg是更稳妥的选择 。它的社区更庞大,设计更偏向于批处理和数据管理。
  • 如果你的场景是基于Spark的增量数据处理,对数据摄取的延迟极其敏感 ,比如需要分钟级甚至秒级将数据更新同步到数据湖供查询, 可以重点评估Hudi ,特别是其Merge-on-Read模式。

Paimon的生态正在快速成长,除了Flink,其与Spark、StarRocks、Doris等的集成也在不断加强。它的定位非常清晰: 在流批一体,特别是Flink流处理优先的架构中,成为数据湖存储层的首选 。它用数据库的思维来管理湖存储,降低了流式数据管理的复杂度,这是它最大的价值所在。

更多推荐