本文介绍如何通过 Serverless StarRocks 查询 Fluss 中的数据。
背景信息
Fluss 支持通过 StarRocks External Catalog 机制,将 Fluss 注册为外部数据源,实现对 Fluss 数据的即席查询和分析。通过 StarRocks,可以:
实时查询 Fluss 中的流式数据,数据延迟为秒级。
查询已分层到 Paimon 的历史数据。
通过 Union Read 合并查询 Fluss 实时数据与 Paimon 历史数据,一条 SQL 即可透明访问冷热全量数据。
前提条件
已创建 Fluss 集群。
已开通 EMR Serverless StarRocks 版实例,且与 Fluss 集群位于同一 VPC 下。详情请参见快速使用存算一体版实例。
Serverless StarRocks 3.3.13-1.2.0 及以上版本。
操作步骤
步骤一:创建 Fluss Catalog
登录 EMR Serverless StarRocks 控制台,在实例列表页面单击目标实例操作列的连接实例。详情请参见通过EMR StarRocks Manager连接StarRocks实例。
在 SQL Editor 中执行以下 SQL,将 Fluss 注册为 External Catalog:
CREATE EXTERNAL CATALOG <catalog_name> PROPERTIES ( 'type' = 'fluss', 'bootstrap.servers' = '<Fluss集群bootstrap地址>', 'fluss.option.client.security.protocol' = 'SASL', 'fluss.option.client.security.sasl.mechanism' = 'PLAIN', 'fluss.option.client.security.sasl.username' = '<Fluss用户名>', 'fluss.option.client.security.sasl.password' = '<Fluss密码>' );参数说明
参数
是否必填
说明
catalog_name
是
Catalog 名称。只能包含字母、数字、下划线,且以字母开头,长度不超过 64 个字符。
type
是
数据源类型,固定为
fluss。bootstrap.servers
是
Fluss 集群的 bootstrap 地址。可在 Fluss 实例详情页获取。
fluss.option.client.security.protocol
是
认证协议。阿里云上默认为
SASL。fluss.option.client.security.sasl.mechanism
是
SASL 认证机制。默认为
PLAIN(账号密码登录)。fluss.option.client.security.sasl.username
是
SASL 认证的账号名。
fluss.option.client.security.sasl.password
是
SASL 认证的账号密码。
Catalog 创建完成后,可以通过三段式命名(
Catalog.Database.Table)直接查询,无需切换上下文。Fluss 的默认数据库名称为fluss:SELECT * FROM <catalog_name>.<database_name>.<table_name>;切换到 Fluss Catalog 后查询:
-- 切换 Catalog 和数据库 SET CATALOG <catalog_name>; USE <db_name>; -- 直接查询目标表 SELECT COUNT(*) FROM <table_name> LIMIT 10;直接切换到对应库查询
USE <catalog_name>.<db_name>;
步骤二:通过后缀指定查询模式
对于已开启湖流一体的 Fluss 表(建表时设置 table.datalake.enabled = 'true'),数据会同时存在于 Fluss 集群(实时热数据)和 Paimon 数据湖(分层历史数据)两个层次。StarRocks 支持通过在表名后添加后缀来指定查询模式,灵活访问不同层次的数据。各查询模式说明如下:
查询模式 | 后缀 | 说明 | 适用场景 |
实时数据查询 |
| 直接查询 Fluss 中尚未归档至数据湖的实时增量数据,数据延迟为秒级。 | 实时监控告警等对延迟敏感的场景。 |
历史数据查询 |
| 查询已分层到 Paimon 的全量历史数据。 | 离线分析、报表统计等场景。 |
合并查询 | 无后缀 | StarRocks 自动合并 Fluss 层(热数据,秒级延迟)与 Paimon 层(冷数据,分钟级延迟),返回完整的最新数据视图。 | 需要同时访问冷热全量数据的场景。 |
以线上销售订单表 order 为例,该表实时记录用户的下单数据(包含订单 ID、订单金额、下单时间、地区等字段),并已开启湖流一体。以下展示三种查询模式的典型业务场景:
-- 场景一:实时监控今日订单量与金额($rt,秒级延迟)
SELECT COUNT(*) AS order_cnt, SUM(order_amount) AS total_amount
FROM `fluss_catalog`.`fluss`.order$rt
WHERE order_time >= CURRENT_DATE;
-- 场景二:计算本周与上周销售额环比增长率($lake,查询 Paimon 历史数据)
SELECT
this_week AS 本周销售额,
last_week AS 上周销售额,
ROUND((this_week - last_week) * 100.0 / NULLIF(last_week, 0), 2) AS 环比增长率
FROM (
SELECT
SUM(CASE WHEN YEARWEEK(order_time) = YEARWEEK(CURRENT_DATE) THEN order_amount ELSE 0 END) AS this_week,
SUM(CASE WHEN YEARWEEK(order_time) = YEARWEEK(CURRENT_DATE) - 1 THEN order_amount ELSE 0 END) AS last_week
FROM `fluss_catalog`.`fluss`.order$lake
) t;
-- 场景三:查询近 7 天各地区的订单总额(无后缀,Union Read 自动合并冷热数据)
SELECT region, SUM(order_amount) AS total_amount
FROM `fluss_catalog`.`fluss`.order
WHERE order_time >= CURRENT_DATE - INTERVAL '7' DAY
GROUP BY region
ORDER BY total_amount DESC;相关操作
查看已创建的 Catalog
-- 查看所有 Catalog
SHOW CATALOGS;
-- 查看指定 Catalog 的创建语句
SHOW CREATE CATALOG <catalog_name>;查看 Fluss 表结构
DESCRIBE <catalog_name>.<database_name>.<table_name>;查看 Fluss 数据库列表
SHOW DATABASES FROM <catalog_name>;删除 Fluss Catalog
删除 Catalog 仅移除 StarRocks 中的映射关系,不会影响 Fluss 中的实际数据。
DROP CATALOG <catalog_name>;