Apache Paimon:基于LSM的流式湖仓存储格式解析与实践
1. 项目概述为什么我们需要一个新的“湖仓”最近几年数据架构领域的热词换了一茬又一茬从数据仓库、数据湖再到湖仓一体技术栈的演进速度让人应接不暇。作为一名在一线折腾了十多年的数据工程师我经历过用Hive做T1批处理的“古典”时代也见证了以Apache Flink为代表的流处理如何重塑实时数仓。但一个核心痛点始终存在我们如何构建一个既能处理海量历史数据湖的灵活性又能支持低延迟实时更新与查询仓的时效性同时还能让流批处理真正统一起来的系统这就是Apache Paimon孵化中要回答的问题。简单来说Paimon是一个流式湖仓存储格式。它不是一个全新的存储引擎而是定义了一套基于LSMLog-Structured Merge-Tree结构的表格式规范。你可以把它理解为一个“超级连接器”它能让Flink、Spark、Hive、Trino这些计算引擎像读写一张普通数据库表一样去读写同一个存储在对象存储如S3、OSS或HDFS上的数据目录并且这张表天然支持流式的读写。我第一次接触Paimon时最直观的感受是它解决了一个非常具体的工程难题流处理的结果如何高效、可靠地沉淀下来并立刻支持点查、批分析甚至流式读取在过去我们可能需要用Flink将实时数据写入Kafka再用一个Flink作业消费Kafka写入Hudi/Iceberg最后才能用Presto去查询。链路长、复杂度高、一致性难保证。Paimon试图将这个过程简化为一步Flink流式写入Paimon表同时任何支持它的引擎都可以实时或批量读取这张表。这听起来像是把数据库的ACID特性和数据湖的开放性、经济性结合了起来。它适合谁呢如果你正在为实时数仓的“最后一公里”——即实时宽表构建、流批一体分析、CDC实时入湖——而头疼或者你对现有湖仓格式在频繁更新、流式读取上的性能不满意那么Paimon值得你花时间深入了解。接下来我会从一个实践者的角度拆解它的核心设计、实操要点以及那些官方文档里不会明说的“坑”。2. 核心设计思路LSM树如何赋能流式湖仓要理解Paimon必须从它的根基——LSM树说起。这不是什么新概念在HBase、RocksDB等系统中早已广泛应用。但Paimon巧妙地将LSM树从单个节点的存储引擎扩展到了分布式文件系统之上从而实现了湖仓的流式能力。2.1 LSM树架构解析写优化与读优化的平衡LSM树的核心思想是“先写日志再后台合并”。在Paimon中这个逻辑被映射到文件系统上写入路径高速写入当数据写入时首先被追加到当前活跃的LSM数据文件中。这个文件是不可变的写满即关闭生成一个新的文件。这个过程完全是顺序I/O写入速度极快这是Paimon能支持高吞吐流式写入的关键。同时所有的写入操作会生成一条记录到WALWrite-Ahead Log文件中用于保证持久性和故障恢复。这就像一家繁忙的餐厅顾客的点单数据先被快速记在便签LSM文件上同时录音笔WAL也在同步录音确保万无一失。合并路径Compaction优化读取随着写入的进行会产生大量小的LSM数据文件。如果直接查询需要扫描大量文件读性能会急剧下降。因此Paimon有一个后台的Compaction进程它会定期将多个小的、可能存在重叠键范围的数据文件合并成更大的、按键排序的文件。合并过程会应用 upsert/delete 语义只保留相同主键的最新记录。这相当于后厨定期将一堆便签整理、誊写到一本整洁的菜单上方便服务员查询引擎快速查找。快照管理一致性视图Paimon通过Snapshot快照来管理数据版本。每次成功的提交一批数据写入完成都会产生一个新的快照。快照记录了当前时刻所有有效数据文件的列表。查询引擎通过指定一个快照ID就能获得一个时间点上一致的数据视图。这种机制天然支持时间旅行查询和增量读取。你可以轻松地查询一小时前的数据全貌或者只读取最近5分钟新增的数据。这种设计带来的直接好处是高吞吐写入追加写模式避开了随机I/O非常适合Flink这类流引擎持续写入。高效更新通过主键定义更新操作被转换为一次新的写入Compaction会负责合并最终状态。流批统一批处理可以读取某个快照的全量数据流处理可以持续读取快照之间的增量数据底层是同一份存储。2.2 与同类技术的差异化定位很多人会问已经有了Apache Iceberg、Apache Hudi和Delta Lake为什么还需要Paimon它们在理念上都属于“湖仓表格式”但侧重点不同。我们可以用一个表格来快速对比它们在核心场景下的差异特性/场景Apache PaimonApache HudiApache IcebergDelta Lake核心设计倾向流式原生读写优化增量处理快速Upsert大规模分析完美模式演化ACID事务与Spark深度集成流式写入原生一流支持写入即流式支持但更强调COW/MOR表类型选择支持但流式读取生态相对较新支持依赖Structured Streaming流式读取核心能力通过快照增量实现支持增量查询支持但需消费元数据日志支持Change Data Feed更新效率高LSM结构擅长处理频繁upsert高MOR表针对更新优化中依赖重写数据文件中依赖重写数据文件批分析性能优合并后文件有序良取决于表类型和压缩优分区与数据组织出色优生态集成深度集成Flink扩展其他引擎与Spark/Flink/Presto集成良好生态最广几乎支持所有引擎与Spark生态绑定最深适用场景实时宽表、CDC实时入湖、流批一体分析近实时数仓、增量ETL超大规模历史数据分析、数据湖治理基于Spark构建的数仓和数据湖从表格可以看出Paimon的杀手锏在于“流式原生”。它从设计之初就为Flink流处理做了深度优化将流处理中的状态管理和数据存储边界模糊化。例如在Flink中你可以将Paimon表直接作为流作业的状态后端的物化视图实现极低延迟的流式更新和查询。这是其他格式目前难以媲美的。注意技术选型没有银弹。如果你的场景是纯批处理历史数据分析Iceberg可能更稳健如果你的团队以Spark技术栈为主Delta Lake上手更快。Paimon的核心优势在于当你已经重度使用Flink构建实时管道并希望简化架构、实现实时查询时它能提供最丝滑的体验。3. 核心细节解析与实操要点理解了设计理念我们深入到实现层面。Paimon表的核心由两部分组成数据文件和元数据。元数据是理解其如何工作的钥匙。3.1 元数据系统如何追踪数据的所有变化Paimon的元数据采用多层结构全部以文件形式存储在对象存储上这保证了其开放性和可审计性。主要包含Snapshot文件snapshot-这是最重要的元数据文件。每次提交生成一个记录了本次提交的版本号、时间戳、所属Schema、包含了哪些数据文件清单Manifest List以及父快照ID。它构成了数据版本链。Manifest List文件一个快照可能引用多个清单文件这个列表文件记录了所有关联的Manifest文件路径。Manifest文件清单文件是数据文件的索引。它记录了多个数据文件的详细元信息包括文件路径、统计信息如每列的最小最大值、空值数、所属分区等。查询时可以根据这些统计信息快速跳过无关文件。Schema文件记录表结构的演变历史。Option文件存储表的配置项如主键、分区键、桶数量、压缩格式等。这种设计的好处是读取操作通常只需要读取最新的Snapshot文件和相关的Manifest文件就能定位到数据无需扫描全量数据目录。写入操作则通过原子性地交换Snapshot文件指针来完成保证了ACID中的原子性和隔离性。3.2 主键表与追加表选择哪一种创建Paimon表时第一个关键决策是选择表类型这直接决定了数据的更新语义和Compaction行为。主键表这是Paimon的推荐和核心表类型。你必须定义主键Primary Key。它的语义类似于数据库表支持根据主键进行Upsert更新插入和Delete。流式写入时相同主键的新数据会覆盖旧数据。适用场景维度表实时更新、CDC同步、实时聚合结果表。例如将MySQL的用户表通过CDC同步到Paimon或者将实时计算的用户画像宽表写入Paimon供查询。注意事项主键字段的选择至关重要。它应该是能唯一标识一行、且更新频率适中的字段。不恰当的主键如更新极其频繁的字段会导致Compaction压力巨大。追加表不定义主键。它只支持Append-Only写入即所有写入的数据都被认为是新的记录即使内容完全相同也不会去重。没有Compaction过程。适用场景事件流水日志、无需更新的事实数据。例如用户点击日志、IoT传感器上报的原始数据流。注意事项由于没有主键无法进行高效的Upsert和点查。如果误将需要更新的数据写入追加表数据会出现重复且无法直接修复。实操心得在实时数仓的维度建模中我通常将维度表设为主键表用于接收CDC数据将事实表设为追加表用于接收事件流。对于需要更新的聚合事实如累计销售额则使用主键表主键设置为“业务日期维度ID”。3.3 分区与分桶数据组织的艺术为了优化查询性能尤其是批处理查询合理的数据组织是必须的。分区与Hive分区类似按照某个时间或类别字段如dt20240501将数据存储到不同目录。分区剪枝能让查询引擎快速跳过无关分区。Paimon支持多级分区如dt20240501/hr10。建议对时间序列数据按天分区是最常见且有效的做法。分区字段不宜过多否则会产生大量小文件。分桶这是Paimon的一个重要特性。在分区内数据可以进一步根据主键或指定字段的哈希值分到固定数量的桶文件中。分桶能保证相同主键的数据一定落在同一个桶文件里这对于主键表的点查和Compaction效率提升巨大。如何设置桶数这是一个经验值。太少会导致单个文件过大影响并行度和Compaction效率太多会产生大量小文件影响HDFS或对象存储性能也增加元数据压力。一个粗略的起点是目标单个桶文件大小在128MB ~ 1GB之间。你可以根据数据总量估算。例如预计一个分区有100GB数据希望每个桶文件约500MB那么桶数可以设置为100GB / 0.5GB 200。配置示例CREATE TABLE user_profile ( user_id BIGINT, name STRING, age INT, city STRING, last_login TIMESTAMP, PRIMARY KEY (user_id) NOT ENFORCED ) PARTITIONED BY (dt STRING) WITH ( bucket 10, -- 每个分区内分10个桶 bucket-key user_id -- 按user_id分桶通常与主键一致 );4. 实操过程从零构建一个实时用户宽表理论说再多不如动手一试。我们以一个经典的场景为例将MySQL的用户基础信息CDC和Kafka中的用户行为日志进行实时关联构建一个实时更新的用户宽表并写入Paimon供即席查询。4.1 环境准备与表定义假设我们已有Flink 1.17环境并已将Paimon Connector的JAR包放入lib/目录。首先在Flink SQL中创建两张源表和一个目标Paimon表。-- 1. 创建MySQL CDC源表捕获用户基础信息变更 CREATE TABLE mysql_user_source ( id BIGINT, name STRING, email STRING, created_at TIMESTAMP(0), updated_at TIMESTAMP(0), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname localhost, port 3306, username flinkuser, password flinkpw, database-name demo, table-name users, server-time-zone Asia/Shanghai ); -- 2. 创建Kafka源表接收用户行为事件 CREATE TABLE kafka_user_behavior ( user_id BIGINT, item_id BIGINT, behavior STRING, -- view, cart, buy event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND ) WITH ( connector kafka, topic user_behavior, properties.bootstrap.servers localhost:9092, scan.startup.mode latest-offset, format json ); -- 3. 创建目标Paimon主键表按天分区按user_id分桶 CREATE TABLE paimon_user_wide_table ( user_id BIGINT, name STRING, email STRING, last_behavior STRING, last_behavior_time TIMESTAMP(3), behavior_count BIGINT, dt STRING, -- 分区字段由事件时间推导 PRIMARY KEY (user_id, dt) NOT ENFORCED -- 联合主键保证每天每个用户一条记录 ) PARTITIONED BY (dt) WITH ( connector paimon, path s3://my-bucket/paimon/user_wide_table, -- 或 hdfs://... auto-create true, bucket 5, -- 根据数据量调整 merge-engine deduplicate -- 默认的去重合并引擎适用于主键表 );关键点解析目标表的主键是(user_id, dt)。这意味着我们的宽表以“用户日期”为粒度每天每个用户只有一条汇总记录。这是实时聚合场景的常见模式。merge-engine选项设置为deduplicate这是主键表的默认引擎保留同一主键下最后一条记录。分区字段dt需要从数据中提取或生成这里我们会在后续SQL中处理。4.2 流式ETL逻辑与写入接下来编写流SQL作业进行关联和聚合并写入Paimon。-- 4. 流处理SQL关联、聚合、写入 INSERT INTO paimon_user_wide_table SELECT u.id as user_id, u.name, u.email, LAST_VALUE(b.behavior) AS last_behavior, -- 获取当天最后一条行为 MAX(b.event_time) AS last_behavior_time, -- 获取当天最后行为时间 COUNT(b.behavior) AS behavior_count, -- 统计当天行为总数 DATE_FORMAT(b.event_time, yyyy-MM-dd) as dt -- 从行为时间生成分区字段 FROM mysql_user_source u JOIN kafka_user_behavior b ON u.id b.user_id WHERE b.behavior IS NOT NULL GROUP BY u.id, u.name, u.email, DATE_FORMAT(b.event_time, yyyy-MM-dd);这个作业的运行逻辑是它是一个持续运行的流作业。实时监听MySQL的users表变更CDC和Kafka的user_behavior主题。将两张流表进行INNER JOIN关联条件是user_id。按用户和事件日期进行分组聚合计算每个用户当天的最后行为、最后行为时间和行为总数。将聚合结果流式地Upsert到Paimon表中。由于主键是(user_id, dt)所以当同一个用户在同一天有新的行为发生时这条记录会被更新last_behavior,last_behavior_time,behavior_count会被新值覆盖。提交这个作业后一个实时更新的用户宽表就在Paimon中自动构建起来了。数据会首先以LSM文件的形式高速写入后台Compaction任务会逐步合并文件优化读取性能。4.3 从多引擎查询数据现在数据已经躺在Paimon里了我们可以用不同的引擎来查询它体验“流式湖仓”的开放性。使用Flink进行流式读取消费增量-- 在另一个Flink SQL任务中流式消费Paimon表的增量变化 SELECT * FROM paimon_user_wide_table /* OPTIONS(scan.modeincremental) */;这个查询会持续输出宽表中发生变更新增或更新的记录可以用于触发下游的实时预警或通知。使用Trino/Presto进行交互式点查-- 在Trino中查询特定用户今天的最新状态 SELECT * FROM paimon.user_wide_table WHERE dt 2024-05-01 AND user_id 123456;得益于主键和分桶这种点查效率非常高几乎可以做到亚秒级响应。使用Spark进行历史批量分析// 在Spark中分析过去一周所有用户的活跃度 val df spark.read.format(paimon).load(s3://my-bucket/paimon/user_wide_table) df.filter($dt 2024-04-24 and $dt 2024-05-01) .groupBy(last_behavior) .count() .show()5. 常见问题与排查技巧实录在实际生产中使用Paimon你肯定会遇到一些挑战。下面是我和团队踩过的一些坑和总结的经验。5.1 小文件问题与Compaction调优问题流作业持续写入即使设置了分桶每个桶内也可能因为频繁提交而产生大量小文件导致查询变慢元数据压力大。根因这是LSM结构的通病。写入频率高、每次写入数据量小微批处理是主要原因。解决方案与调优调整写入提交间隔在Flink的Table配置中增加sink.parallelism并调大checkpoint间隔或sink.buffer-flush.interval让每次提交包含更多数据从而生成更大的数据文件。-- 在WITH参数中调整 WITH ( ... sink.parallelism 4, sink.buffer-flush.max-rows 10000, -- 缓冲1万行刷一次 sink.buffer-flush.interval 30s -- 或30秒刷一次 )调整Compaction策略Paimon的Compaction是自动触发的但可以调整其策略。compaction.max.file-num设置一个桶内触发Compaction的最大文件数默认是50。可以适当调小如30让合并更频繁避免文件堆积。compaction.target-file-size合并后目标文件大小默认128MB。如果你的存储和计算资源充足可以适当调大如256MB或512MB减少文件总数。注意更频繁的Compaction会消耗更多的CPU和I/O资源需要在文件数量和资源消耗之间取得平衡。建议在测试环境进行压测观察不同参数下的表现。5.2 流读延迟与消费位点管理问题使用scan.modeincremental进行流式读取时发现数据延迟很高或者重启任务后从很旧的位置开始消费。根因Paimon的流读本质是追踪快照Snapshot的增量。消费位点保存在状态中。如果任务失败且没有配置合适的快照保留策略或者位点信息丢失就可能出现问题。排查与解决检查快照过期策略Paimon默认会保留所有快照这会导致元数据目录膨胀。但如果你配置了snapshot.time-retained如snapshot.time-retained 1 h来定期清理过期快照那么流读任务如果挂掉超过1小时就可能找不到对应的快照而失败。对于重要的流读任务建议设置较长的保留时间或使用snapshot.num-retained.min和snapshot.num-retained.max来控制数量。明确指定起始快照在流读SQL中可以使用/* OPTIONS(scan.snapshot-idxxx) */来指定从某个具体的快照ID开始消费。你可以通过Paimon的系统表sys.snapshots查询现有快照。监控消费延迟可以通过查询sys.snapshots表对比当前最新快照的ID和流读任务消费的快照ID来计算延迟。5.3 分区与主键设计陷阱问题查询性能不佳点查依然很慢。排查检查分区字段是否在查询条件中如果你的查询总是带dt条件但表没有按dt分区那么每次查询都会扫描全表数据。确保高频查询条件作为分区字段。检查分桶键是否合理分桶键默认与主键一致。如果主键是(user_id, dt)那么分桶的哈希值是基于这两个字段计算的。这能保证同一个用户同一天的数据在一个桶内。不要选择高基数列如时间戳单独作为分桶键这会导致数据严重倾斜所有数据可能都集中到少数几个桶。避免“数据倾斜”如果某个分区或某个桶的数据量远大于其他部分会导致Compaction和查询长尾。对于分区可以考虑使用更细粒度如小时对于分桶如果某个键值特别多如“未知用户”可以考虑在ETL阶段进行预处理。5.4 写入失败与数据一致性问题Flink作业写入Paimon失败重启后发现有数据重复或丢失。行动清单确保启用Checkpoint这是Flink保证精确一次Exactly-Once语义的基础。Paimon Sink依赖Checkpoint来提交事务生成快照。检查WAL日志Paimon的WAL用于在故障恢复时重放未提交的数据。确保WAL目录默认在表路径下的_wal目录有足够的空间和写入权限。理解“批模式”写入在非流式作业如Spark批作业中写入Paimon默认是“批模式”即整个作业完成后一次性提交。如果作业中途失败不会有部分数据写入。但在流式场景下需要依赖Checkpoint。善用sys.options表Paimon提供了系统表来查看和修改表的配置。当遇到问题时可以查询SELECT * FROM sys.options$table_name来确认当前生效的配置。最后我个人在实际操作中的体会是Paimon将流处理的思维深深植入到了数据存储层这种“流原生”的特性让实时数仓的链路变得异常简洁。它尤其适合作为Flink流作业的物化结果表和CDC入湖的统一出口。但它的成熟度仍在快速演进中社区和生态是你在选型时需要重点评估的因素。对于已经拥有成熟Hudi或Iceberg平台的公司引入Paimon可能会增加技术栈的复杂度但对于从零开始构建、且以Flink为核心实时计算引擎的团队Paimon提供了一个极具吸引力的、端到端的流式湖仓解决方案。开始使用时建议从一个非核心的实时场景入手逐步摸清其特性和最佳实践再向核心业务推广。