【数据库】Fluss 是什么?如何使用 Fluss 构建流式湖仓一体?
Fluss是一款具备亚秒级低延迟特性的流式存储系统,支持高效的实时数据读写。借助Lakehouse Storage,Fluss在Lakehouse架构之上构建了实时流数据服务,从而实现了数据流与数据湖仓的无缝集成与统一管理。
一、湖仓概述
1. 传统湖仓架构的挑战
数据湖仓(Lakehouse)作为一种创新的开放架构,融合了数据湖的灵活扩展性和经济高效性,以及数据仓库的稳定可靠性和高性能优势。在这一架构中,Apache Iceberg、Apache Paimon、Apache Hudi 和 Delta Lake 等领先的数据湖格式发挥了重要作用,它们通过统一的平台,实现了数据存储、可靠性与分析能力之间的高效协同与平衡。
然而,传统湖仓架构面临着实时性与分析能力平衡的挑战:
(1) 实时性与分析效率的矛盾
- 如果要求低延迟,就需要频繁写入和提交,这会产生大量小型Parquet文件,导致读取效率低下
- 如果要求读取效率,就需要积累数据直到能写入大型Parquet文件,但这会引入更高的延迟
(2) 数据新鲜度限制
即使在最佳使用条件下,这些数据湖格式通常也只能在分钟级粒度内实现数据新鲜度。

二、流式湖仓一体化
1. 流与湖的统一
Fluss是一种支持亚秒级低延迟流式读写的流式存储系统。通过Lakehouse Storage,Fluss在Lakehouse之上提供实时流数据服务,实现了数据流和数据湖仓的统一。这不仅为数据湖仓带来了低延迟,还为数据流增加了强大的分析能力。
为了构建流式湖仓,Fluss维护了一个分层服务,该服务将实时数据从Fluss集群压缩到存储在Lakehouse Storage中的数据湖格式。Fluss集群中的数据(流式Arrow格式)针对低延迟读写进行了优化,而Lakehouse中的压缩数据(带压缩的Parquet格式)针对强大的分析和长期数据存储进行了优化。因此,Fluss集群中的数据作为实时数据层,保留了具有亚秒级新鲜度的数据;而Lakehouse中的数据作为历史数据层,保留了具有分钟级新鲜度的数据。
2. 共享数据与共享元数据
流式湖仓的核心理念是流和湖仓之间共享数据和共享元数据,避免数据重复和元数据不一致。它提供了以下强大功能:
- 统一元数据:Fluss为流和湖仓中的数据提供统一的表元数据。用户只需处理一个表,但可以访问实时流数据、历史数据或它们的联合。
- 联合读取:计算引擎对表执行查询时将读取实时流数据和湖仓数据的联合。目前,只有Flink支持联合读取,但更多引擎正在规划中。
- 实时湖仓:联合读取帮助湖仓从近实时分析发展到真正的实时分析,使企业能够从实时数据中获得更有价值的洞察。
- 分析流:联合读取帮助数据流具备强大的分析能力,这减少了开发流应用程序的复杂性,简化了调试,并允许立即访问实时数据洞察。
- 连接到湖仓生态系统:Fluss在将数据压缩到Lakehouse时保持表元数据与数据湖目录同步,这允许Spark、StarRocks、Flink、Trino等外部引擎通过连接到数据湖目录直接读取数据。
3. 数据湖集成
(1) Paimon集成
Apache Paimon创新地结合了湖格式和LSM结构,将高效更新引入湖架构。
要将Fluss与Paimon集成,必须启用lakehouse存储并将Paimon配置为lakehouse存储。具体步骤如下:
① 配置Lakehouse存储
首先,在server.yaml中配置lakehouse存储:
复制
datalake.format: paimon
# Paimon目录配置,假设使用Filesystem目录
datalake.paimon.metastore: filesystem
datalake.paimon.warehouse: /tmp/paimon_data_warehouse
② 启动数据湖分层服务
然后,启动数据湖分层服务,将Fluss的数据压缩到lakehouse存储:
复制
# 切换到Fluss目录
cd $FLUSS_HOME
# 启动分层服务,假设rest端点是localhost:8081
./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081
③ 为表启用Lakehouse存储
要为表启用lakehouse存储,必须在创建表时使用选项'table.datalake.enabled' = 'true':
复制
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');
当在Fluss中创建或修改带有选项'table.datalake.enabled' = 'true'的表时,Fluss将创建一个具有相同表路径的相应Paimon表。Paimon表的模式与Fluss表的模式相同,只是在最后附加了两个额外的列__offset和__timestamp。这两列用于帮助Fluss客户端以流式方式消费Paimon中的数据,例如按偏移量/时间戳查找等。
然后,数据湖分层服务会持续将数据从Fluss压缩到Paimon。对于主键表,它还将生成Paimon格式的变更日志,使您能够以Paimon方式流式消费它。
(2) Iceberg集成规划
目前,Fluss支持Paimon作为Lakehouse存储,更多种类的数据湖格式正在规划中。
Fluss社区正在积极努力增强流和湖仓统一功能,重点关注以下关键领域:
① 扩展联合读取生态系统
目前,联合读取功能已与Apache Flink集成,实现了实时和历史数据的无缝查询。未来,社区计划扩展此功能以支持其他查询引擎,如Apache Spark和StarRocks,进一步扩大其生态系统兼容性和采用。
② 多样化湖存储格式
目前,Fluss支持Apache Paimon作为其主要湖存储。为了满足多样化的用户需求,社区计划添加对更多湖格式的支持,包括Apache Iceberg和Apache Hudi,从而提供更大的灵活性和与更广泛的湖仓生态系统的互操作性。
4. 联合读取(Union Read)
(1) 实时数据与历史数据的无缝结合
对于带有选项'table.datalake.enabled' = 'true'的表,有两部分数据:保留在Fluss中的数据和已经在Paimon中的数据。现在,您有两种表视图:一种是具有分钟级延迟的Paimon数据视图,一种是联合Fluss和Paimon数据的完整数据视图,具有秒级延迟。
Flink使您能够决定选择哪种视图:
- 仅Paimon意味着更好的分析性能,但数据新鲜度较差
- 结合Fluss和Paimon意味着更好的数据新鲜度,但分析性能下降
① 仅读取Paimon中的数据
要指定读取Paimon中的数据,必须使用$lake后缀指定表,以下SQL显示了如何执行此操作:
复制
-- 假设我们有一个名为`orders`的表
-- 从paimon读取
SELECT COUNT(*) FROM orders$lake;
-- 我们还可以查询系统表
SELECT * FROM orders$lake$snapshots;
当在查询中使用$lake后缀指定表时,它就像一个普通的Paimon表,因此它继承了Paimon表的所有能力。您可以享受Flink在Paimon上支持/优化的所有功能,如查询系统表、时间旅行等。
② 联合读取Fluss和Paimon中的数据
要指定读取联合Fluss和Paimon的完整数据,只需像查询普通表一样查询它,无需任何后缀或其他内容,以下SQL显示了如何执行此操作:
复制
-- 查询将联合Fluss和Paimon的数据
SELECT SUM(order_count) as total_orders FROM ads_nation_purchase_power;
查询可能看起来比只查询Paimon中的数据慢,但它查询了完整数据,这意味着更好的数据新鲜度。您可以多次运行查询,由于数据持续写入表中,每次运行都应该得到不同的结果。
(2) 查询引擎支持
① Flink支持
Fluss向Flink用户公开统一的API,允许他们选择是使用联合读取还是仅在Lakehouse上进行读取,使用以下SQL:
SELECT * FROM orders
这将读取orders表的完整数据,Flink将联合读取Fluss和Lakehouse中的数据。如果用户只需要读取数据湖上的数据,可以在要读取的表后添加$lake后缀。SQL如下:
-- 分析查询
SELECT COUNT(*), MAX(t), SUM(amount)
FROM orders$lake
-- 查询系统表
SELECT * FROM orders$lake$snapshots
5. 实时湖仓与分析流
(1) 为湖仓带来真正的实时性
传统的数据湖仓架构通常只能提供分钟级的数据新鲜度,这对于许多实时分析场景来说是不够的。Fluss的流式湖仓一体化架构通过联合读取(Union Read)机制,为湖仓带来了真正的实时性:
- 秒级数据新鲜度:通过将Fluss的实时流数据与Paimon的历史数据联合起来,查询可以访问最新写入的数据,实现秒级数据新鲜度。
- 实时分析能力增强:企业可以在保持湖仓强大分析能力的同时,获得实时数据洞察,使决策更加及时和准确。
- 无需数据复制:不需要将数据复制到单独的实时系统中,减少了存储成本和数据一致性问题。
(2) 为数据流增加强大的分析能力
传统的流处理系统通常缺乏强大的分析能力,而Fluss通过与湖仓的集成,为数据流增加了这些能力:
- 复杂分析查询支持:流数据可以与历史数据结合,支持更复杂的分析查询,如聚合、窗口分析等。
- 开发简化:减少了开发流应用程序的复杂性,简化了调试过程。
- 生态系统集成:可以利用现有的湖仓生态系统工具和查询引擎,如Spark、StarRocks、Trino等。
6. 流式湖仓一体化实施流程
下面是实施Fluss流式湖仓一体化的详细步骤和流程图:

(1) 准备环境
确保已安装Fluss和Flink,并且两者都正常运行。Fluss需要一个运行中的Flink集群来执行数据湖分层服务。
(2) 配置Lakehouse存储
在Fluss的server.yaml配置文件中设置Lakehouse存储:
# 指定数据湖格式为Paimon
datalake.format: paimon
# Paimon目录配置
datalake.paimon.metastore: filesystem
datalake.paimon.warehouse: /path/to/paimon/warehouse
(3) 启动数据湖分层服务
数据湖分层服务是一个Flink作业,负责将数据从Fluss压缩到Paimon:
复制
# 切换到Fluss安装目录
cd $FLUSS_HOME
# 启动分层服务,指定Flink REST端点
./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081
您还可以设置其他Flink配置参数,例如:
# 设置检查点间隔为10秒
./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081 -D flink.execution.checkpointing.interval=10s
# 仅同步特定数据库中的表
./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081 -D database=fluss_\\w+
(4) 创建启用Lakehouse的表
使用Flink SQL创建启用了Lakehouse的表,通过设置'table.datalake.enabled' = 'true'选项:
复制
-- 创建一个启用了Lakehouse的日志表
CREATE TABLE log_orders (
order_id BIGINT,
customer_id BIGINT,
order_time TIMESTAMP,
amount DECIMAL(10, 2),
status STRING
) WITH (
'bucket.num' = '4',
'table.datalake.enabled' = 'true'
);
-- 创建一个启用了Lakehouse的主键表
CREATE TABLE pk_orders (
order_id BIGINT,
customer_id BIGINT,
order_time TIMESTAMP,
amount DECIMAL(10, 2),
status STRING,
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'bucket.num' = '4',
'table.datalake.enabled' = 'true'
);
当创建启用了Lakehouse的表时,Fluss会自动创建一个对应的Paimon表,并在表末尾添加两个额外的列__offset和__timestamp,用于支持流式消费。
(5) 数据写入
使用标准的Flink SQL向表中写入数据:
复制
-- 插入数据到日志表
INSERT INTO log_orders
VALUES (1001, 101, TIMESTAMP '2025-05-18 10:30:00', 199.99, 'PENDING');
-- 插入数据到主键表
INSERT INTO pk_orders
VALUES (1001, 101, TIMESTAMP '2025-05-18 10:30:00', 199.99, 'PENDING');
-- 或者从其他表插入数据
INSERT INTO pk_orders
SELECT order_id, customer_id, order_time, amount, status
FROM source_table;
数据写入后,会首先存储在Fluss的实时层中,然后由数据湖分层服务定期压缩到Paimon中。对于主键表,还会生成Paimon格式的变更日志。
(6) 数据查询
① 联合读取(实时+历史数据)
要查询完整的数据(Fluss中的实时数据和Paimon中的历史数据),直接查询表名,无需任何后缀:
复制
-- 联合读取实时和历史数据
SELECT COUNT(*) FROM pk_orders;
-- 流式查询,会持续获取最新数据
SELECT * FROM pk_orders WHERE amount > 100;
联合读取提供秒级数据新鲜度,但可能会比仅查询Paimon数据慢一些。由于数据持续写入,多次运行相同的查询可能会得到不同的结果。
② 仅读取Lakehouse数据
要仅查询Paimon中的历史数据,在表名后添加$lake后缀:
-- 仅查询Paimon中的历史数据
SELECT COUNT(*) FROM pk_orders$lake;
-- 查询Paimon系统表
SELECT * FROM pk_orders$lake$snapshots;
-- 使用Paimon的时间旅行功能
SELECT * FROM pk_orders$lake FOR TIMESTAMP AS OF '2025-05-17 00:00:00';
仅读取Paimon数据提供更好的分析性能,但数据新鲜度较差(分钟级)。您可以使用Paimon的所有功能,如查询系统表、时间旅行等。
7. 流式湖仓一体化架构图

三、实际应用场景示例
场景1:实时仪表盘
一个电子商务平台需要一个实时仪表盘,显示最新的销售数据和趋势。
实现方式:
- 创建启用Lakehouse的订单表
- 使用Flink SQL进行联合读取,获取最新的销售数据
- 将结果推送到仪表盘应用
-- 创建启用Lakehouse的订单表
CREATE TABLE sales_orders (
order_id BIGINT,
product_id BIGINT,
customer_id BIGINT,
order_time TIMESTAMP,
amount DECIMAL(10, 2),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'bucket.num' = '8',
'table.datalake.enabled' = 'true'
);
-- 实时计算最近1小时的销售额
SELECT
TUMBLE_START(order_time, INTERVAL '5' MINUTE) AS window_start,
SUM(amount) AS total_sales
FROM sales_orders
WHERE order_time > CURRENT_TIMESTAMP - INTERVAL '1' HOUR
GROUP BY TUMBLE(order_time, INTERVAL '5' MINUTE);
场景2:历史数据分析与实时数据结合
一个金融机构需要分析客户的历史交易模式,并将其与实时交易数据结合,以检测潜在的欺诈行为。
实现方式:
- 创建启用Lakehouse的交易表
- 使用Paimon查询历史交易模式
- 使用联合读取将历史模式与实时交易结合分析
-- 创建启用Lakehouse的交易表
CREATE TABLE transactions (
transaction_id BIGINT,
account_id BIGINT,
transaction_time TIMESTAMP,
amount DECIMAL(15, 2),
merchant_id BIGINT,
location STRING,
PRIMARY KEY (transaction_id) NOT ENFORCED
) WITH (
'bucket.num' = '16',
'table.datalake.enabled' = 'true'
);
-- 查询历史交易模式(仅Paimon数据)
SELECT
account_id,
AVG(amount) AS avg_amount,
STDDEV(amount) AS stddev_amount
FROM transactions$lake
WHERE transaction_time BETWEEN TIMESTAMP '2025-04-01 00:00:00' AND TIMESTAMP '2025-05-01 00:00:00'
GROUP BY account_id;
-- 将历史模式与实时交易结合分析(联合读取)
WITH historical_patterns AS (
SELECT
account_id,
AVG(amount) AS avg_amount,
STDDEV(amount) AS stddev_amount
FROM transactions$lake
WHERE transaction_time BETWEEN TIMESTAMP '2025-04-01 00:00:00' AND TIMESTAMP '2025-05-01 00:00:00'
GROUP BY account_id
)
SELECT
t.transaction_id,
t.account_id,
t.transaction_time,
t.amount,
t.merchant_id,
t.location,
h.avg_amount,
h.stddev_amount,
CASE
WHEN ABS(t.amount - h.avg_amount) > 3 * h.stddev_amount THEN 'SUSPICIOUS'
ELSE 'NORMAL'
END AS risk_level
FROM transactions t
JOIN historical_patterns h ON t.account_id = h.account_id
WHERE t.transaction_time > CURRENT_TIMESTAMP - INTERVAL '1' HOUR;
场景3:实时数据湖分层架构
一个大型零售企业需要构建一个数据湖分层架构(Bronze、Silver、Gold),同时保持各层的数据新鲜度。
实现方式:
- 创建各层的启用Lakehouse的表
- 使用Flink SQL将数据从一层转换到下一层
- 利用联合读取确保各层都具有秒级数据新鲜度

复制
-- Bronze层:创建原始销售数据表
CREATE TABLE sales_raw (
event_time TIMESTAMP,
store_id BIGINT,
product_id BIGINT,
quantity INT,
amount DECIMAL(10, 2),
payment_method STRING,
source STRING
) WITH (
'bucket.num' = '16',
'table.datalake.enabled' = 'true'
);
-- Silver层:创建清洗后的销售数据表
CREATE TABLE sales_cleaned (
event_time TIMESTAMP,
store_id BIGINT,
product_id BIGINT,
quantity INT,
amount DECIMAL(10, 2),
payment_method STRING,
source STRING,
PRIMARY KEY (store_id, product_id, event_time) NOT ENFORCED
) WITH (
'bucket.num' = '16',
'table.datalake.enabled' = 'true'
);
-- Gold层:创建聚合的销售指标表
CREATE TABLE sales_aggregated (
window_start TIMESTAMP,
window_end TIMESTAMP,
store_id BIGINT,
product_id BIGINT,
total_quantity BIGINT,
total_amount DECIMAL(15, 2),
PRIMARY KEY (window_start, window_end, store_id, product_id) NOT ENFORCED
) WITH (
'bucket.num' = '8',
'table.datalake.enabled' = 'true'
);
-- Silver层ETL:清洗和转换数据
INSERT INTO sales_cleaned
SELECT
event_time,
store_id,
product_id,
CASE WHEN quantity <= 0 THEN 1 ELSE quantity END AS quantity,
amount,
payment_method,
source
FROM sales_raw
WHERE amount > 0;
-- Gold层ETL:聚合数据
INSERT INTO sales_aggregated
SELECT
TUMBLE_START(event_time, INTERVAL '1' HOUR) AS window_start,
TUMBLE_END(event_time, INTERVAL '1' HOUR) AS window_end,
store_id,
product_id,
SUM(quantity) AS total_quantity,
SUM(amount) AS total_amount
FROM sales_cleaned
GROUP BY TUMBLE(event_time, INTERVAL '1' HOUR), store_id, product_id;
四、流式湖仓一体化的技术优势
1. 数据分布一致性
Fluss和Paimon之间的数据分布是严格对齐的。Fluss支持分区表和桶,其分桶算法与Paimon完全一致。这确保了给定的数据始终分配到两个系统中的相同桶,创建了Fluss桶和Paimon桶之间的一对一对应关系。
这种数据分布的一致性提供了两个显著优势:
(1) 消除分层过程中的Shuffle开销
当将数据从Fluss分层到Paimon格式时:
- Fluss桶(例如bucket1)可以直接分层到相应的Paimon桶(bucket1)
- 无需读取Fluss桶中的数据,计算每条数据属于哪个Paimon桶,然后将其写入适当的Paimon桶
通过绕过这个中间重新分配步骤,架构避免了昂贵的shuffle开销,显著提高了压缩效率。
(2) 防止数据不一致
通过在Fluss和Paimon中使用相同的分桶算法,确保了数据一致性。该算法计算每条数据的桶分配如下:
bucket_id = hash(row) % bucket_num
2. 更高效的数据追踪
在Fluss架构中,历史数据存储在Lakehouse存储中,而实时数据保留在Fluss中。在流式读取期间,这种架构实现了历史和实时数据访问的无缝结合:
- 历史数据访问:Fluss直接从Lakehouse存储检索历史数据,利用其固有优势,如:
- 高效的过滤下推:使查询引擎能够在存储层应用过滤条件,减少读取的数据量并提高性能
- 列裁剪:允许仅检索必要的列,优化数据传输和查询效率
- 高压缩比:在保持快速检索速度的同时最小化存储开销
- 实时数据访问:Fluss同时从自己的存储中读取最新的实时数据,确保毫秒级的新鲜度
3. 一致的数据新鲜度
在构建数据仓库时,通常按层组织和管理数据,如Bronze、Silver和Gold层。当数据流经这些层时,维护数据新鲜度变得至关重要。
当Paimon作为每一层的唯一存储解决方案时,数据可见性取决于Flink检查点间隔,这会引入累积延迟:
- 给定层的变更日志仅在Flink检查点完成后才可见
- 当这个变更日志传播到后续层时,数据新鲜度延迟随着每个检查点间隔而增加
例如,使用1分钟的Flink检查点间隔:
- Bronze层经历1分钟延迟
- Silver层增加另一个1分钟延迟,总计2分钟
- Gold层再增加1分钟延迟,累积为3分钟延迟
而使用Fluss和Paimon,我们获得:
- 即时数据可见性:Fluss中的数据在摄入后立即可见,无需等待Flink检查点完成。变更日志立即传输到下一层
- 一致的数据新鲜度:所有层的数据新鲜度是一致的,以秒为单位,消除了累积延迟
4. 完整实施示例:电子商务实时分析平台
下面是一个完整的电子商务实时分析平台实施示例,展示了如何使用Fluss的流式湖仓一体化架构构建端到端解决方案。
步骤1:环境准备
首先,确保已安装并配置好Fluss和Flink环境:
# 下载并解压Fluss
wget https://github.com/alibaba/fluss/releases/download/v0.6.0/fluss-0.6.0.tgz
tar -xzf fluss-0.6.0.tgz
cd fluss-0.6.0
# 启动本地Fluss集群
./bin/start-local.sh
# 启动本地Flink集群(用于数据湖分层服务)
cd /path/to/flink
./bin/start-cluster.sh
步骤2:配置Lakehouse存储
在Fluss的conf/server.yaml中配置Paimon作为Lakehouse存储:
datalake.format: paimon
datalake.paimon.metastore: filesystem
datalake.paimon.warehouse: /path/to/paimon/warehouse
步骤3:启动数据湖分层服务
启动数据湖分层服务,将数据从Fluss压缩到Paimon:
cd /path/to/fluss
./bin/lakehouse.sh -D flink.rest.address=localhost -D flink.rest.port=8081 -D flink.execution.checkpointing.interval=10s
步骤4:创建数据模型
使用Flink SQL创建电子商务分析所需的表结构:
-- 创建订单表(主键表)
CREATE TABLE orders (
order_id BIGINT,
customer_id BIGINT,
order_time TIMESTAMP,
total_amount DECIMAL(10, 2),
status STRING,
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'bucket.num' = '8',
'table.datalake.enabled' = 'true'
);
-- 创建订单明细表(日志表)
CREATE TABLE order_items (
item_id BIGINT,
order_id BIGINT,
product_id BIGINT,
quantity INT,
unit_price DECIMAL(10, 2),
discount DECIMAL(5, 2)
) WITH (
'bucket.num' = '8',
'table.datalake.enabled' = 'true'
);
-- 创建产品表(主键表)
CREATE TABLE products (
product_id BIGINT,
name STRING,
category STRING,
brand STRING,
price DECIMAL(10, 2),
PRIMARY KEY (product_id) NOT ENFORCED
) WITH (
'bucket.num' = '4',
'table.datalake.enabled' = 'true'
);
-- 创建客户表(主键表)
CREATE TABLE customers (
customer_id BIGINT,
name STRING,
email STRING,
registration_time TIMESTAMP,
vip_level INT,
PRIMARY KEY (customer_id) NOT ENFORCED
) WITH (
'bucket.num' = '4',
'table.datalake.enabled' = 'true'
);
-- 创建实时销售指标表(主键表)
CREATE TABLE sales_metrics (
window_start TIMESTAMP,
window_end TIMESTAMP,
category STRING,
total_orders BIGINT,
total_sales DECIMAL(15, 2),
PRIMARY KEY (window_start, window_end, category) NOT ENFORCED
) WITH (
'bucket.num' = '4',
'table.datalake.enabled' = 'true'
);
步骤5:数据处理流程
使用Flink SQL实现数据处理流程:
-- 1. 计算实时销售指标(每5分钟更新一次)
INSERT INTO sales_metrics
SELECT
TUMBLE_START(o.order_time, INTERVAL '5' MINUTE) AS window_start,
TUMBLE_END(o.order_time, INTERVAL '5' MINUTE) AS window_end,
p.category,
COUNT(DISTINCT o.order_id) AS total_orders,
SUM(i.quantity * i.unit_price * (1 - i.discount/100)) AS total_sales
FROM orders o
JOIN order_items i ON o.order_id = i.order_id
JOIN products p ON i.product_id = p.product_id
WHERE o.status = 'COMPLETED'
GROUP BY
TUMBLE(o.order_time, INTERVAL '5' MINUTE),
p.category;
-- 2. 创建VIP客户实时推荐视图
CREATE VIEW vip_recommendations AS
SELECT
c.customer_id,
c.name,
p.category,
p.brand,
COUNT(*) AS purchase_count
FROM customers c
JOIN orders o ON c.customer_id = o.customer_id
JOIN order_items i ON o.order_id = i.order_id
JOIN products p ON i.product_id = p.product_id
WHERE
c.vip_level >= 3 AND
o.order_time > CURRENT_TIMESTAMP - INTERVAL '30' DAY
GROUP BY
c.customer_id, c.name, p.category, p.brand
HAVING COUNT(*) >= 2;
步骤6:实时仪表盘查询
使用联合读取查询实时销售指标:
-- 实时销售趋势(联合读取)
SELECT
window_start,
category,
total_sales
FROM sales_metrics
WHERE
window_start >= CURRENT_TIMESTAMP - INTERVAL '24' HOUR
ORDER BY
window_start DESC, total_sales DESC;
步骤7:历史分析查询
使用Paimon查询历史数据分析:
-- 历史销售分析(仅Paimon数据)
SELECT
p.category,
MONTH(o.order_time) AS month,
SUM(i.quantity * i.unit_price * (1 - i.discount/100)) AS total_sales
FROM orders$lake o
JOIN order_items$lake i ON o.order_id = i.order_id
JOIN products$lake p ON i.product_id = p.product_id
WHERE
o.order_time >= TIMESTAMP '2025-01-01 00:00:00' AND
o.order_time < TIMESTAMP '2025-05-01 00:00:00'
GROUP BY
p.category, MONTH(o.order_time)
ORDER BY
p.category, month;
这种查询只访问Paimon中的历史数据,提供更好的查询性能,适合大规模历史数据分析。
5. 流式湖仓一体化的关键技术点
(1) 统一元数据管理
在传统架构中,流存储系统(如Kafka)和湖仓存储解决方案(如Paimon)作为独立实体运行,各自维护自己的元数据。这给计算引擎(如Flink)带来了两个主要挑战:
双重目录:用户需要创建和管理两个独立的目录——一个用于流存储,另一个用于湖存储
手动切换:访问数据需要手动在目录之间切换,以确定是查询流存储还是湖存储,导致操作复杂性和效率低下
而在Fluss中,虽然Fluss和Paimon仍然维护独立的元数据,但它们向计算引擎(如Flink)公开统一的目录和单一的表抽象。这种统一方法提供了几个关键优势:
- 简化数据访问:用户可以通过单一目录无缝访问湖仓存储(Paimon)和流存储(Fluss),无需管理或在独立目录之间切换
- 集成查询:统一表抽象允许直接访问Fluss中的实时数据和Paimon中的历史数据
- 操作效率:通过提供一致的接口,该架构降低了操作复杂性,使用户更容易在单一工作流中处理实时和历史数据
(2) 数据分布对齐
Fluss和Paimon之间的数据分布是严格对齐的,这是通过使用相同的分桶算法实现的:
bucket_id = hash(row) % bucket_num
这种对齐确保了:
- 分层效率:Fluss桶可以直接分层到对应的Paimon桶,无需重新分配数据
- 数据一致性:相同的数据在两个系统中分配到相同的桶,防止数据不一致
(3) 检查点间隔与数据新鲜度
在传统的Paimon架构中,数据新鲜度受Flink检查点间隔的限制。例如,使用1分钟的检查点间隔时:

而在Fluss流式湖仓架构中,数据在写入后立即可见:

6. 流式湖仓一体化的最佳实践
(1) 表设计最佳实践
- 合理设置桶数量:根据数据量和查询模式设置适当的桶数量,通常为2的幂次方(4、8、16等)
- 选择合适的主键:对于主键表,选择具有良好分布特性的主键,避免热点
- 考虑分区策略:对于大型表,使用分区来提高查询性能和管理数据生命周期
(2) 数据湖分层服务配置
- 检查点间隔:根据实时性需求和系统负载调整检查点间隔,通常在10-60秒之间
- 资源分配:为数据湖分层服务分配足够的资源,确保能够及时处理数据
- 监控:定期监控分层服务的性能和延迟,确保数据及时压缩到Paimon
( 3)查询优化
①选择合适的查询模式
- 对于需要最新数据的查询,使用联合读取
- 对于大规模历史数据分析,使用仅Paimon查询
②过滤下推:在查询中尽可能早地应用过滤条件,利用Paimon的过滤下推能力
③列裁剪:只选择需要的列,减少数据传输和处理量
更多推荐
所有评论(0)