Serverless StarRocks访问Fluss(Beta)

更新时间:
复制 MD 格式

本文介绍如何通过 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

  1. 登录 EMR Serverless StarRocks 控制台,在实例列表页面单击目标实例操作列的连接实例。详情请参见通过EMR StarRocks Manager连接StarRocks实例。

  2. 在 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 认证的账号密码。

  3. 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 支持通过在表名后添加后缀来指定查询模式,灵活访问不同层次的数据。各查询模式说明如下:

查询模式

后缀

说明

适用场景

实时数据查询

$rt

直接查询 Fluss 中尚未归档至数据湖的实时增量数据,数据延迟为秒级。

实时监控告警等对延迟敏感的场景。

历史数据查询

$lake

查询已分层到 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>;