Fluss初窥五StreamingLakeHouse
LakeHouse(湖仓一体)
LakeHouse(湖仓一体)是一种新型的开放式存储计算架构,打通了数据仓库和数据湖,将数据湖的灵活性、可扩展性和成本效益与数据仓库的可靠性和高性能及管理能力相结合。底层支持多种数据类型并存,能实现数据间的相互共享,上层可以通过统一封装的接口进行访问,可同时支持实时查询和分析,为企业进行数据治理带来了更多的便利性。
常见的数据湖格式,如 Apache Iceberg、Apache Paimon、Apache Hudi 和 Delta Lake,在湖仓一体架构中发挥着关键作用,促进了单一统一平台内数据存储、可靠性和分析能力之间的和谐平衡。
Lakehouse 作为一种新型架构,能够有效地满足数据管理和分析的复杂需求。但是,它们很难满足实时分析的场景,因为实时分析需要亚秒级的数据新鲜度,而这些数据新鲜度受到其实施的限制。使用这些数据湖格式,您将陷入矛盾的情况:
- 如果你需要低延迟,那么就需要频繁地写入和提交,这意味着会产生许多小的 Parquet 文件。对于必须处理大量小文件的读取操作,这变得效率低下。
- 如果您需要读取效率,那么你会积累数据,直到可以写入大型 Parquet 文件,但这会引入更高的延迟。
总而言之,这些数据湖格式在最佳使用条件下,数据新鲜度最多只能达到分钟级粒度。
Streaming Lakehouse(流式湖仓)
Fluss 是一种支持流式读取和写入,具有亚秒级低延迟的流式存储。与 Lakehouse Storage 结合,Fluss 通过在 Lakehouse 之上提供实时流数据,统一了数据流和数据湖。这不仅为数据湖带来了低延迟,还为数据流增加了强大的分析功能。
为了构建 Streaming Lakehouse(流式湖仓),Fluss 维护了一个压缩服务,将 Fluss 集群中的实时数据压缩到 Lakehouse 存储中。Fluss 集群中的数据(流式 Arrow 格式)针对低延迟读写进行了写优化,Lakehouse 中的压缩数据(带压缩的 Parquet 格式)针对强大的分析进行了读优化,并针对存储长期数据进行了空间优化。
因此,Fluss 集群中的数据充当实时数据层,保留数天的数据,并保持亚秒级的新鲜度,而湖仓中的数据充当历史数据层,保留数月的数据,并保持分钟级的新鲜度。

流式湖仓的核心思想是在流和湖仓之间共享数据和元数据,避免数据重复和元数据不一致。它提供的一些强大功能包括:
- Unified Metadata(统一元数据): Fluss 为 Stream 和 Lakehouse 中的数据提供统一的表元数据。因此,用户只需处理一个表,但可以访问实时流数据、历史数据或它们的并集。
- Union Reads(联合读取): 计算引擎在表上执行查询,将读取实时流数据和湖仓数据的并集。目前,只有 Flink 支持联合读取,但更多引擎正在规划中。
- 实时湖仓: 联合读取帮助湖仓从近似实时分析发展到真正实时分析。使得企业能够从实时数据中获得更有价值的洞察。
- 流式分析: 联合读取帮助数据流具有强大的分析功能。这降低了开发流式应用程序的复杂性,简化了调试,并允许立即访问实时数据洞察。
- 连接 Lakehouse 生态系统: Fluss 将表元数据与数据湖目录同步,同时将数据压缩到 Lakehouse 中。这使得外部引擎如 Spark、StarRocks、Flink、Trino 可以直接通过连接到数据湖目录来读取数据。
湖仓存储
湖仓存储集群配置
首先,必须在 server.yaml 中配置湖屋存储。以派蒙为例,您必须配置以下设置:
yaml
datalake.tiered.storage: paimon
# the catalog config about Paimon, assuming using Filesystem catalog
paimon.catalog.type: filesystem
paimon.catalog.warehouse: /tmp/paimon_data_warehouse
启动数据湖分层服务
然后,必须启动数据湖分层服务以将 Fluss 的数据压缩到湖仓存储。要启动数据湖分层服务,您必须有一个正在运行的 Flink 集群,因为 Fluss 目前仅支持 Flink 作为分层服务后端。
bash
# change directory to Fluss
cd $FLUSS_HOME
# start the tiering service, assuming rest endpoint is localhost:8081
./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081
flink.rest.address和flink.rest.port是 Flink 集群的 REST 端点,您可能需要根据您的 Flink 集群配置进行更改。- 数据湖分层服务实际上是一个 Flink 作业,您可以在启动数据湖分层服务时通过
-D参数设置 Flink 配置。例如,如果您想将检查点间隔设置为 10 秒,可以使用以下命令启动数据湖分层服务:
bash
./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081 -D flink.execution.checkpointing.interval=10s
启用每张表的湖仓存储
要启用表的数据湖存储,必须使用选项 table.datalake.enabled = true 创建表。
对 Paimon 的支持
Apache Paimon 创新性地结合了湖格式和 LSM 结构,为湖架构带来了高效的更新。要集成 Fluss 与 Paimon,您必须启用湖屋存储并配置 Paimon 为湖屋存储。当在 Fluss 中创建或修改带有选项 table.datalake.enabled = true 的表时,Fluss 将创建一个具有相同表路径的相应 Paimon 表。Paimon 表的架构与 Fluss 表的架构相同,除了在最后附加了两个额外的列 __offset 和 __timestamp。这两个列用于帮助 Fluss 客户端以流式方式消费 Paimon 中的数据,例如通过偏移量/时间戳等进行查找。然后数据湖分层服务持续压缩 Fluss 到 Paimon 的数据。对于主键表,它还会以 Paimon 格式生成变更日志,使您能够以 Paimon 的方式流式消费。
读取表
由 Flink 读取
对于选项 table.datalake.enabled = true 的表格,数据分为两部分:Fluss 中保留的数据和已在 Paimon 中的数据。现在有两种查看表格的方式:一种是 Paimon 数据视图,具有分钟级延迟,另一种是 Fluss 和 Paimon 数据的完整数据联合视图,最新数据在秒级延迟内。Flink 允许您决定选择哪种视图:
- 仅 Paimon 意味着更好的分析性能,但数据新鲜度更差
- Fluss 和 Paimon 的结合意味着更好的数据新鲜度,但分析性能会下降
仅读取 Paimon 中的数据
指向 Paimon 中的读取数据,您必须指定带有 $lake 后缀的表,以下 SQL 展示了如何操作:
sql
-- assume we have a table named `orders`
-- read from paimon
SELECT COUNT(*) FROM orders$lake;
-- we can also query the system tables
SELECT * FROM orders$lake$snapshots;
当在查询中使用带有 $lake 后缀的表时,它就像一个普通的 Paimon 表一样,因此继承了 Paimon 表的所有能力。您可以享受 Flink 查询在 Paimon 上支持/优化的所有功能,如查询系统表、时间旅行等。有关 Paimon 的 SQL 查询的更多信息,请参考 Paimon 官网。
联合读取 Fluss 和 Paimon 中的数据
要指向读取 Fluss 和 Paimon 联合数据的完整数据,您只需将其查询为一个普通表,无需任何后缀或其他,以下 SQL 展示了如何操作:
sql
-- query will union data of Fluss and Paimon
SELECT SUM(order_count) as total_orders FROM ads_nation_purchase_power;
查询可能看起来比仅在 Paimon 中查询数据要慢,但它查询的是全部数据,这意味着数据更新更好。可以多次运行查询,每次运行都应该得到不同的结果,因为数据是持续写入表中的。
由其他引擎读取
以 StarRocks 作为读取数据的引擎,首先,为 StarRocks 创建一个 Paimon 目录:
sql
CREATE EXTERNAL CATALOG paimon_catalog
PROPERTIES
(
"type" = "paimon",
"paimon.catalog.type" = "filesystem",
"paimon.catalog.warehouse" = "/tmp/paimon_data_warehouse"
);
注意:配置值 paimon.catalog.type 和 paimon.catalog.warehouse 应与您在 server.yaml 中将 Paimon 配置为 Fluss 湖仓存储的方式相同。然后,可以使用 StarRocks 查询 orders 表:
sql
-- the table is in database `fluss`
SELECT COUNT(*) FROM paimon_catalog.fluss.orders;
-- query the system tables, to know the snapshots of the table
SELECT * FROM paimon_catalog.fluss.enriched_orders$snapshots;
FlinkSQL + Fluss -> Paimon Demo
启动 Lakehouse 分层服务
要集成 Apache Paimon,您需要启动 Lakehouse Tiering Service。打开一个新的终端,导航到 fluss-quickstart-flink 目录,并在该目录中执行以下命令以启动服务:
bash
./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081 -D flink.execution.checkpointing.interval=10s
流式传输到 Fluss - 数据湖启用表
默认情况下,创建的表具有禁用数据湖集成,这意味着湖仓分层服务不会将表的数据分层到数据湖。要启用湖仓功能作为表的分层存储解决方案,您必须使用配置选项 table.datalake.enabled = true 创建表。返回到 SQL client 并执行以下 SQL 语句以创建启用数据湖集成的表:
sql
CREATE TABLE datalake_enriched_orders (
`order_key` BIGINT,
`cust_key` INT NOT NULL,
`total_price` DECIMAL(15, 2),
`order_date` DATE,
`order_priority` STRING,
`clerk` STRING,
`cust_name` STRING,
`cust_phone` STRING,
`cust_acctbal` DECIMAL(15, 2),
`cust_mktsegment` STRING,
`nation_name` STRING,
PRIMARY KEY (`order_key`) NOT ENFORCED
) WITH ('table.datalake.enabled' = 'true');
将流式数据写入启用数据湖的表中,datalake_enriched_orders:
sql
-- switch to streaming mode
SET 'execution.runtime-mode' = 'streaming';
INSERT INTO datalake_enriched_orders
SELECT o.order_key,
o.cust_key,
o.total_price,
o.order_date,
o.order_priority,
o.clerk,
c.name,
c.phone,
c.acctbal,
c.mktsegment,
n.name
FROM fluss_order o
LEFT JOIN fluss_customer FOR SYSTEM_TIME AS OF `o`.`ptime` AS `c`
ON o.cust_key = c.cust_key
LEFT JOIN fluss_nation FOR SYSTEM_TIME AS OF `o`.`ptime` AS `n`
ON c.nation_key = n.nation_key;
实时分析 Fluss 数据湖
datalake_enriched_orders 的数据存储在 Fluss(实时数据)和 Paimon(历史数据)中。在查询 datalake_enriched_orders 表时,Fluss 使用了一个联合操作,该操作将 Fluss 和 Paimon 的数据组合在一起,以提供一个完整的结果集——结合了实时和历史数据。如果只想查询存储在 Paimon 中的数据——提供高性能访问,无需合并数据的开销——您可以通过添加 $lake 表。此方法还启用所有 Flink Paimon 表源优化和功能,包括系统表如 datalake_enriched_orders$snapshots。
从 Paimon 直接查询快照
sql
-- switch to batch mode
SET 'execution.runtime-mode' = 'batch';
-- to query snapshots in paimon
SELECT snapshot_id, total_record_count FROM datalake_enriched_orders$lake$snapshots;
+-------------+--------------------+
| snapshot_id | total_record_count |
+-------------+--------------------+
| 1 | 650 |
+-------------+--------------------+
注意:在查询快照之前,请确保等待检查点(约 30 秒)完成,否则结果将为空。
使用 SQL 语句对 Paimon 数据进行统计分析
sql
-- to sum prices of all orders in paimon
SELECT sum(total_price) as sum_price FROM datalake_enriched_orders$lake;
直接查询表
联合读取返回的是 Fluss + Paimon 的最新数据,满足亚秒级数据新鲜度的结果。
sql
-- to sum prices of all orders in fluss and paimon
SELECT sum(total_price) as sum_price FROM datalake_enriched_orders;
+------------+
| sum_price |
+------------+
| 1777908.36 |
+------------+
使用以下命令查看存储在 Paimon 中的文件
bash
tree /tmp/paimon/fluss.db
/tmp/paimon/fluss.db
└── datalake_enriched_orders
├── bucket-0
│ ├── changelog-aef1810f-85b2-4eba-8eb8-9b136dec5bdb-0.orc
│ └── data-aef1810f-85b2-4eba-8eb8-9b136dec5bdb-1.orc
├── manifest
│ ├── manifest-aaa007e1-81a2-40b3-ba1f-9df4528bc402-0
│ ├── manifest-aaa007e1-81a2-40b3-ba1f-9df4528bc402-1
│ ├── manifest-list-ceb77e1f-7d17-4160-9e1f-f334918c6e0d-0
│ ├── manifest-list-ceb77e1f-7d17-4160-9e1f-f334918c6e0d-1
│ └── manifest-list-ceb77e1f-7d17-4160-9e1f-f334918c6e0d-2
├── schema
│ └── schema-0
└── snapshot
├── EARLIEST
├── LATEST
└── snapshot-1
大数据相关学习资料
大数据相关学习资料、大数据项目、湖仓一体、架构师必知必会、数据中台建设方法论...共有 1400 多份文档资料,另专为星球成员整理了一份比较详细的语雀知识库合计 190 万字和飞书文档资料。(内容太多,仅展示部分内容...),欢迎大家踊跃加入星球,您将获得:
- 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
- 免费答疑解惑、交流技术。
- 面试指导、模拟面试。
- 各类 PDF 文档下载、星球代码下载。
- 提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导。
另外说明加入星球后支持三天无理由退款,不满意无条件随时退。需要资料请加微信:D1435221412

end
