Apache Paimon:基于LSM树与Flink深度集成的流批一体数据湖存储
1. 从“数据仓库”到“数据湖”为什么我们需要Paimon如果你在过去几年里处理过海量数据尤其是流式数据那么“数据湖”这个概念对你来说一定不陌生。传统的数仓模式比如Hive在处理T1的批处理任务时表现尚可但一旦面对实时数据流就显得力不从心。数据延迟、复杂的ETL链路、批流两套系统带来的维护成本都成了数据工程师的痛点。而数据湖的初衷就是提供一个统一的存储层既能存原始数据又能支持批处理和流处理让数据“随到随用”。然而理想很丰满现实很骨感。早期的数据湖方案比如基于HDFS和Hive表虽然解决了存储统一的问题但在实时更新、事务一致性、流批一体查询性能上依然存在巨大挑战。简单来说它们更像一个“数据沼泽”——数据扔进去容易想高效、一致地拿出来用却困难重重。正是在这样的背景下以Apache Iceberg、Apache Hudi和Delta Lake为代表的“湖仓一体”Lakehouse架构应运而生它们试图在数据湖的灵活性和数据仓库的性能与管理能力之间架起一座桥。那么Paimon是什么它正是在这个赛道上的一个新选手而且是带着鲜明的流处理基因出生的。Paimon原名Flink Table Store最初是Apache Flink社区内部孵化的项目目标非常明确为Flink流处理引擎打造一个高性能、低延迟的流批一体存储层。你可以把它理解成Flink的“原生”数据湖格式。它不是为了取代谁而是为了填补一个特定的空白当你的技术栈核心是Flink你需要一个能完美配合Flink流式读写、Exactly-Once语义、以及实时更新的存储系统时Paimon可能就是那个最“趁手”的工具。它的核心价值在于将流处理中“状态”的概念与数据湖“表”的概念深度融合。在Flink流作业中状态State是计算正确性的基石但它通常存储在RocksDB这样的嵌入式KV存储里难以直接查询和归档。Paimon则把这个状态“外化”成了一个可查询、可回溯、可批量处理的数据湖表。这意味着流处理作业的中间结果或最终结果可以直接沉淀为一张湖表供下游的即时查询OLAP、批量分析或另一个流作业消费真正实现了流批的存储与计算统一。2. Paimon的架构全景一张表的“三层楼”设计要理解Paimon必须从它的物理存储架构入手。Paimon的表数据在文件系统如HDFS、S3、OSS上以一种清晰的分层结构组织。我习惯把它比喻成一栋“三层楼”的建筑每一层都有其特定的职责。2.1 底层Snapshot快照—— 表的“版本目录”这是Paimon架构中最核心的元数据层。快照Snapshot定义了表在某个时间点的完整状态。你可以把它想象成Git中的一个Commit或者相机的一张照片它捕获了表在某个瞬间的所有数据。是什么一个快照文件snapshot-xxx本质上是一个JSON文件它记录了本次提交的元信息比如快照ID、所属的分区、包含了哪些数据文件Manifest List、以及Schema信息等。干什么用它是实现时间旅行Time Travel和增量读取的基础。通过指定快照ID或时间戳你可以查询到表在历史上任意时刻的样子。这对于数据回溯、审计、修正错误数据回滚到某个快照至关重要。如何工作每次向Paimon表写入数据一个Checkpoint或一次手动提交都会生成一个新的快照。Paimon会维护一个快照链最新的快照指向当前表的数据。这种设计使得读操作无论是批读还是流读都能获得一个一致性的视图不会读到正在写入的“脏数据”。2.2 中间层Manifest清单—— 数据的“文件索引”如果说快照是目录那么清单就是目录里的索引页。一个快照会引用一个或多个清单文件Manifest List而每个清单文件则记录了更详细的数据文件信息。是什么清单文件manifest-xxx也是一个元数据文件它存储了一个数据文件Data File列表。每个条目会包含数据文件的路径、统计信息如每列的最大最小值用于谓词下推加速查询、文件大小、记录条数等。干什么用它充当了查询优化器的“导航图”。当执行一个带有过滤条件的查询时例如WHERE user_id 123Paimon的读取器可以快速扫描清单文件中的统计信息直接跳过那些根本不可能包含目标数据的数据文件极大地提升了查询效率。这是数据湖格式相比原始文件存储如Parquet/ORC直接放在HDFS上的核心优势之一。清单列表Manifest List为了管理海量数据文件Paimon引入了清单列表。一个快照引用一个清单列表文件该文件再指向多个具体的清单文件。这是一种分层索引避免单个清单文件过大。2.3 顶层Data File数据文件—— 真正的“货物仓库”这里存放着用户的实际数据以列式存储格式默认是ORC也支持Parquet存储。格式ORC格式因其高效的压缩和读取性能被选为Paimon的默认数据文件格式。它非常适合分析型查询。组织方式数据文件通常按分区和桶Bucket进行组织。分区如dt2023-10-01是常见的粗粒度数据划分方式。桶则是Paimon实现高效更新和避免小文件的关键机制我们稍后会详细讲。合并Compaction在流式写入场景下会持续产生小的数据文件。Paimon后台有合并任务负责将多个小文件合并成大文件并清理无效数据被后续更新或删除覆盖的数据这个过程对于维持查询性能至关重要。这三层结构共同保证了Paimon表的数据一致性、高效的查询性能以及丰富的数据管理功能时间旅行、增量读取。理解了这个“三层楼”模型你就掌握了Paimon物理存储的骨架。3. 核心原理深度拆解Paimon如何实现流式更新与高效查询了解了静态架构我们再来看看Paimon动态工作的核心原理。这主要集中在它如何应对流处理中最棘手的两个问题持续更新和高效点查。3.1 LSM树结构流式更新的引擎Paimon表的核心存储引擎借鉴了LSM-TreeLog-Structured Merge-Tree的思想。LSM树是很多现代数据库如Google Bigtable、Cassandra、RocksDB处理高频写入的基石。它的核心思想是“先写日志再后台合并”将随机写转换为顺序写从而大幅提升写入吞吐。在Paimon的语境下LSM树是如何工作的呢写入阶段顺序追加当一条数据写入Paimon表时它并不是直接去修改已有的ORC数据文件那会是昂贵的随机IO。相反Paimon会先将这条写入记录可能包含插入、更新、删除以顺序追加的方式写入一个特殊的文件——变更日志文件Changelog File。你可以把它理解成LSM树的“MemTable”刷盘后的形态。这个过程非常快。读取阶段合并视图当用户查询表时Paimon的读取器需要提供一个统一的数据视图。它会同时读取基础数据文件Sorted Runs和最新的变更日志文件在内存中按照主键进行合并将最终结果返回给用户。对于一条记录如果在变更日志中有新的版本就会覆盖基础数据文件中的旧版本。合并阶段Compaction变更日志文件会不断累积如果每次查询都合并大量小文件效率会很低。因此Paimon后台会定期触发合并Compaction任务。这个任务将多个变更日志文件以及它们覆盖的旧基础数据文件合并生成新的、有序的基础数据文件并生成新的快照。同时旧的、被覆盖的数据文件会被标记为可删除由过期快照清理任务处理。这个过程就是LSM树的“归并排序”过程。这种设计带来的核心优势高吞吐写入写入永远是顺序追加不受数据量和更新模式影响。高效的更新和删除更新和删除操作被转换为新增一条带标记的记录在合并时进行处理原生支持。流式读取Flink CDC等工具捕获的数据库变更日志可以非常自然地作为流直接写入Paimon形成端到端的流式数仓。注意Paimon提供了两种合并策略lookup和full-compaction。lookup合并只合并变更日志适用于点查频繁的场景full-compaction会全量合并基础数据和变更日志能提供最好的读取性能但资源消耗更大。需要根据业务场景权衡选择。3.2 主键表与桶Bucket组织数据的艺术并非所有Paimon表都采用LSM树。Paimon根据表是否有主键分为两类主键表Primary Key Table这是Paimon的“完全体”支持更新和删除操作。它必须定义主键并且必须指定分桶Bucket键通常就是主键或主键的子集。数据根据桶键的哈希值被分配到固定数量的桶中。每个桶内部数据文件按主键排序。这种“分区内分桶桶内有序”的结构是Paimon实现高效点查和范围查询的关键。为什么分桶分桶将数据打散到多个目录中实现了数据的并行读写。更重要的是它将主键的查询范围缩小到了一个桶内。当你要根据主键查询一条记录时系统只需计算其桶号然后去那个桶目录下查找大大减少了需要扫描的数据量。桶内有序在每个桶内部数据文件无论是基础文件还是变更日志都是按照主键排序的。这使得在合并时可以使用高效的归并排序并且在点查时可以使用二分查找等优化手段。仅追加表Append-Only Table没有定义主键只能插入数据不能更新或删除。它的实现更简单类似于传统的分区Hive表数据直接以ORC文件形式追加写入分区目录。适用于日志、事件流等不需要更新的场景。桶数量的选择是一个重要的调优点桶数太少每个桶的数据量过大导致合并压力大点查性能下降并行度也不够。桶数太多会产生大量小文件管理开销大也可能影响某些查询的性能。经验法则通常建议桶的数量与处理数据的并发度如Flink作业的并行度相匹配或成倍数关系并且确保每个桶最终的数据量在1GB到几个GB的合理范围内。可以在建表时通过bucket N来指定。3.3 索引与查询优化如何快速找到数据Paimon的查询性能很大程度上依赖于其元数据提供的索引信息。分区剪枝Partition Pruning这是最基础的优化。如果查询条件中包含了分区字段的过滤如dt ‘2023-10-01’Paimon的规划器会直接跳过所有其他分区的数据只读取目标分区。桶过滤Bucket Filtering对于主键表如果查询条件包含了完整的桶键通常是主键的等值条件Paimon可以精确计算出记录所在的桶直接读取那个桶的数据实现点查Point Lookup。谓词下推Predicate Pushdown这是清单文件统计信息发挥作用的地方。Paimon的读取器在扫描清单文件时会利用每个数据文件记录的列级统计信息最小值、最大值、空值数等。例如查询WHERE age 30如果一个数据文件的age列最大值是25那么这个文件会被直接跳过根本不会被打开。数据跳过Data Skipping除了文件级的统计信息ORC/Parquet文件格式内部还有行组Row Group或条带Stripe级别的统计信息。Paimon可以进一步利用这些信息在打开文件后跳过不可能包含目标数据的行组。这些优化手段层层递进使得Paimon在面对海量数据时依然能提供可接受的查询延迟特别是对于点查和带有过滤条件的分析查询。4. Paimon与Flink的深度集成流批一体的实践Paimon脱胎于Flink其与Flink的集成深度是其他数据湖格式难以比拟的。这种集成不是简单的“能读能写”而是从API到运行时机制的深度融合。4.1 作为流式的Source与Sink这是Paimon最自然的用法。在Flink SQL或DataStream API中你可以轻松地将一个Paimon表定义为源表Source或目标表Sink。流式写入Sink你可以将Kafka数据流、CDC变更流、或经过处理的Flink数据流直接写入一张Paimon主键表。Paimon Sink会以Checkpoint为周期将接收到的数据插入、更新、删除提交为新的快照。这相当于拥有了一个实时更新的、可查询的流状态。流式读取Source有界流Bounded读取表在某个快照下的全部数据用于批处理或初始化。无界流Unbounded这是Paimon的杀手锏功能——流式读取增量数据。你可以通过SELECT * FROM table_name /* OPTIONS(‘scan.mode’‘latest’) */这样的语法或指定scan.mode参数让作业持续监控Paimon表每当有新的快照生成就读取自上一个快照以来变化的数据Changelog。这使得Paimon表可以作为一个流式管道连接上下游作业实现复杂的流式处理链路。4.2 流式维表关联Temporal Join在实时数仓中流事实表关联维表是一个高频场景。传统的方案是查询外部数据库如Redis、HBase带来延迟和压力。Paimon提供了一个优雅的解决方案流式维表关联。你可以将维度数据写入一张Paimon主键表。当流事实表需要关联维表时Flink运行时可以直接将Paimon维表以时态表Temporal Table的形式接入。关联时Flink会根据事实记录的事件时间去查找Paimon维表在该时间点的快照版本从而实现精准的“时间旅行”关联。这比查询外部数据库更高效且能保证维度数据变更的历史准确性。4.3 物化视图与流式聚合Paimon可以作为流式聚合结果的存储。例如你可以创建一个Flink作业实时计算每分钟的销售额并将结果持续写入一张Paimon表。这张结果表本身就是一张物化视图。下游的BI工具或即席查询可以直接查询这张Paimon表获得亚秒级延迟的聚合结果。Paimon的更新能力保证了聚合结果的实时性。4.4 一致性保证得益于与Flink Checkpoint机制的深度集成Paimon的写入可以完美支持Flink的精确一次Exactly-Once语义。在发生故障恢复时可以确保数据既不丢失也不重复。同时快照机制提供了读一致性查询者总是能看到一个完整的快照不会读到部分提交的数据。5. 实战指南从零开始使用Paimon理论说了这么多我们来点实际的。假设我们要构建一个简单的用户行为实时统计场景。5.1 环境准备与表创建首先你需要一个Flink环境1.17版本推荐并下载Paimon的jar包。这里以Flink SQL Client为例。-- 1. 创建Catalog指定Paimon表存储的Warehouse路径例如HDFS或S3 CREATE CATALOG paimon_catalog WITH ( ‘type’‘paimon’, ‘warehouse’‘hdfs:///path/to/warehouse’ ); USE CATALOG paimon_catalog; -- 2. 创建一张主键表用于存储用户明细支持更新 -- 这里我们按dt分区按user_id分桶5个桶并定义user_id为主键 CREATE TABLE user_behavior ( dt STRING, user_id BIGINT, item_id BIGINT, behavior STRING, ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL ‘5’ SECOND, PRIMARY KEY (dt, user_id) NOT ENFORCED -- 主键分区字段必须包含在主键内 ) PARTITIONED BY (dt) WITH ( ‘bucket’ ‘5’, -- 分桶数建议与并行度相关 ‘changelog-producer’ ‘full-compaction’ -- 指定合并生产者保证流读时有完整的changelog ); -- 3. 创建一张仅追加表用于存储聚合结果每分钟用户行为计数 CREATE TABLE behavior_agg ( window_start TIMESTAMP(3), user_id BIGINT, behavior STRING, cnt BIGINT ) WITH ( ‘bucket’ ‘-1’ -- 追加表可以不分桶或设置为-1 );5.2 流式写入与读取-- 1. 流式写入假设有一个Kafka数据源 kafka_source INSERT INTO user_behavior SELECT DATE_FORMAT(ts, ‘yyyy-MM-dd’) as dt, user_id, item_id, behavior, ts FROM kafka_source; -- 2. 流式增量读取启动一个作业实时消费user_behavior表的变更 -- 这个作业会一直运行每当user_behavior表有新数据提交这里就会读到 SELECT * FROM user_behavior /* OPTIONS(‘scan.mode’‘latest’) */; -- 3. 批处理读取查询今天的所有数据 SELECT * FROM user_behavior WHERE dt ‘2023-10-27’; -- 4. 时间旅行查询昨天下午3点整的表状态 SELECT * FROM user_behavior /* OPTIONS(‘scan.snapshot-id’‘12345’) */; -- 或 SELECT * FROM user_behavior FOR SYSTEM_TIME AS OF TIMESTAMP ‘2023-10-26 15:00:00’;5.3 核心参数调优与避坑经验在实际使用中以下几个参数的配置对性能和稳定性影响巨大snapshot.time-retained快照保留时间。默认可能只保留1小时。务必根据业务需要调大例如‘snapshot.time-retained’ ‘7d’否则你将无法进行一周前的时间旅行查询过期快照对应的数据文件也可能被删除。changelog-producer决定如何为流式读取生成变更日志。可选none,input,lookup,full-compaction。full-compaction在合并时生成完整的变更日志能提供最高效的流读性能但会消耗更多计算资源进行合并。生产环境流读链路的源表推荐使用此模式。lookup在流读时通过查找数据文件生成变更日志节省写入端资源但可能增加读取延迟。需根据资源瓶颈权衡。compaction.early-max.file-num触发合并的早期文件数量阈值。如果流写入速度很快会产生大量小文件。调低此值如从50调到10可以让合并更频繁避免小文件过多但会增加IO压力。需要监控文件数量来调整。write-buffer-size和‘compaction.max-size-amplification-percent’写入缓冲和合并策略参数影响写入性能和空间放大。通常默认值即可在写入压力极大时可适当调大写入缓冲区。小文件问题这是所有数据湖格式的共性问题。除了调整合并参数确保写入作业的并行度与桶数匹配避免单个并行度写入多个桶导致文件碎片化。定期监控HDFS上的文件数量和大小是关键。Schema EvolutionPaimon支持添加列、删除列等Schema变更。但修改列类型或重命名列可能不被支持或行为与预期不符。进行Schema变更前务必在测试环境充分验证并备份数据。Paimon作为Flink生态的原生数据湖存储在流批一体、实时更新和高效查询方面展现出了独特的优势。它尤其适合已经以Flink为核心流处理引擎的数据平台用于构建实时数仓、流式物化视图和流批融合的分析场景。当然它并非万能在超大规模纯批处理、或需要与Spark、Trino等引擎进行深度、高性能交互的场景下可能还需要评估与Iceberg、Hudi的成熟度差异。但无论如何Paimon的出现为流处理领域的数据存储提供了一个强有力的、高度集成的选择。