湖流一体表通过统一的元数据层,实现对实时增量数据(Fluss)与历史全量数据(Paimon)的透明访问。可根据业务需求,灵活选择实时流读、离线批读或混合读取模式。
支持的查询引擎
湖流一体表的访问兼容性取决于访问方式。
访问方式 | 说明 | 支持的引擎 |
Fluss 原生访问(含 Union Read) | 充分利用湖流一体的混合读取能力,实现流批统一。 | Flink |
底层 Paimon 访问( | 直接访问底层 Paimon 格式存储,支持多引擎分析。 | Flink、Spark、Trino、StarRocks 等所有兼容 Paimon 的引擎 |
通过 Flink 访问
实时消费
在流处理场景中,通常需要监听数据的最新变更。通过配置 Flink SQL 的启动模式,可以指定不同的消费模式。
-- 实时监听 orders 表的变更,仅读取最新到达的数据
SELECT *
FROM fluss_catalog.fluss.orders
/*+ OPTIONS('scan.startup.mode' = 'latest-offset') */;更多启动模式配置(如最早位点、特定时间戳等),请参见消费模式中的消费模式说明。
批量查询
Fluss 服务会自动将数据实时同步至底层的 Paimon 格式存储。对于需要分析全量历史数据或访问特定快照的离线任务,可直接查询底层的 Paimon 表。
在 Fluss Catalog 中,通过在表名后附加 $lake 后缀,即可直接访问对应的 Paimon 表。
-- 批模式:直接查询底层 Paimon 存储的全量数据
SELECT * FROM fluss_catalog.fluss.orders$lake;
-- 批模式:查询底层 Paimon 存储的快照(Snapshots)元数据
SELECT * FROM fluss_catalog.fluss.orders$lake$snapshots;$lake 表本质上是一张标准的 Paimon 表,因此继承了 Paimon 的所有特性,包括时间旅行(Time Travel)和对多引擎的支持。
关联查询
关联查询(Union Read)提供"湖流融合"的统一视图。系统自动合并 Fluss 层(热数据,秒级延迟)与 Paimon 层(冷数据,分钟级延迟),对外呈现为一张逻辑完整的表,无需关心数据的物理存储位置。
在传统的流批分离架构中,若要同时获取实时与离线数据,需自行搭建数据管道完成数据迁移、ETL 聚合、主键去重等流程,开发维护成本高,且难以保证数据一致性。Union Read 在存储引擎层面原生融合流与湖,一条 SQL 即可透明访问冷热全量数据,无需关心底层数据的物理分片与合并逻辑。
-- 自动合并读取冷热数据
SELECT * FROM fluss_catalog.fluss.orders;Flink 引擎的执行逻辑如下。
执行模式 | 逻辑 |
流模式 | 作业启动时,首先全量读取 Paimon 中的历史数据,随即无缝切换至 Fluss 消费毫秒级的实时增量数据。 |
批模式 | 读取 Paimon 中的基量数据与 Fluss 中的未归档增量数据,并进行合并,返回最新的全局快照。 |
对于主键表,Flink 在批模式下必须对两部分数据执行按主键合并去重。相比日志表的读取,此过程会增加一定的计算耗时。
更多使用Flink访问详情,请参见Fluss Connector。
通过 StarRocks 访问
完成 Fluss Catalog 的创建和配置后(详见引擎对接 - StarRocks),可通过以下方式查询湖流一体表中的数据。
查询湖上数据
StarRocks 通过 $lake 后缀直接访问底层 Paimon 存储的全量历史数据。
-- 查询底层 Paimon 存储的全量数据
SELECT * FROM <catalog_name>.<database_name>.<table_name>$lake;查询流上数据
StarRocks 通过 $rt 后缀直接查询 Fluss 中尚未归档至数据湖的实时增量数据。
-- 查询 Fluss 层的实时增量数据
SELECT * FROM <catalog_name>.<database_name>.<table_name>$rt;查询全部数据(Union Read)
不加任何后缀时,StarRocks 自动合并 Fluss 层(热数据)与 Paimon 层(冷数据),返回完整的最新数据视图。
-- 自动合并冷热数据,返回全量最新数据
SELECT * FROM <catalog_name>.<database_name>.<table_name>;