改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

Iceberg与StarRocks深度集成实战指南

ckckck2025年1月10日3 浏览

1: 配置 StarRocks 访问 Iceberg Catalog

使用 SQL 客户端连接到 StarRocks

MySQL CLI

您可以从 Docker 环境或您的本机运行此客户端。

bash 复制代码
docker compose exec starrocks-fe 
mysql -P 9030 -h 127.0.0.1 -u root --prompt="StarRocks > "

如果您使用 StarRocks 容器中的 MySQL Client,需要从包含 docker-compose.yml 文件的路径运行以下命令。

bash 复制代码
docker compose exec starrocks-fe 
mysql -P 9030 -h 127.0.0.1 -u root --prompt="StarRocks > "

2: 创建 External Catalog

可以通过创建 External Catalog 将 StarRocks 连接至您的数据湖。以下示例基于以上 Iceberg 数据源创建 External Catalog。

sql 复制代码
CREATE EXTERNAL CATALOG 'iceberg'
PROPERTIES
(
  "type"="iceberg",
  "iceberg.catalog.type"="rest",
  "iceberg.catalog.uri"="http://iceberg-rest:8181",
  "iceberg.catalog.warehouse"="warehouse",
  "aws.s3.access_key"="admin",
  "aws.s3.secret_key"="password",
  "aws.s3.endpoint"="http://minio:9000",
  "aws.s3.enable_path_style_access"="true",
  "client.factory"="com.starrocks.connector.iceberg.IcebergAwsClientFactory"
);

创建成功后,运行以下命令查看创建的 Catalog。

sql 复制代码
SHOW CATALOGS;

其中 default_catalog 为 StarRocks 的 Internal Catalog,用于存储内部数据。设置当前使用的 Catalog 为 iceberg

sql 复制代码
SET CATALOG iceberg;

查看 iceberg 中的数据库。

sql 复制代码
SHOW DATABASES;

3: 使用 StarRocks 查询 Iceberg

查询接单时间

以下语句查询出租车接单时间,仅返回前十行数据。

sql 复制代码
SELECT lpep_pickup_datetime FROM greentaxis LIMIT 10;

查询接单高峰时期

以下查询按每小时聚合行程数据,计算每小时接单的数量。

sql 复制代码
SELECT COUNT(*) AS trips,
       hour(lpep_pickup_datetime) AS hour_of_day
FROM greentaxis
GROUP BY hour_of_day
ORDER BY trips DESC;

4: 创建 Iceberg 表

同 StarRocks 内部数据库一致,如果您拥有 Iceberg 数据库的 CREATE TABLE 权限,那么可以使用 CREATE TABLECREATE TABLE AS SELECT (CTAS) 在该 Iceberg 数据库下创建表。本功能自 3.1 版本起开始支持。

sql 复制代码
CREATE TABLE [IF NOT EXISTS] [database.]table_name
(column_definition1[, column_definition2, ...
partition_column_definition1,partition_column_definition2...])
[partition_desc]
[PROPERTIES ("key" = "value", ...)]
[AS SELECT query]

创建非分区表 unpartition_tbl

包含 idscore 两列,如下所示:

sql 复制代码
CREATE TABLE unpartition_tbl
(
    id int,
    score double
);

创建分区表 partition_tbl_1

包含 actioniddt 三列,并定义 iddt 为分区列,如下所示:

sql 复制代码
CREATE TABLE partition_tbl_1
(
    action varchar(20),
    id int,
    dt date
)
PARTITION BY (id,dt);

查询原表 partition_tbl_1 的数据

并根据查询结果创建分区表 partition_tbl_2,定义 iddtpartition_tbl_2 的分区列:

sql 复制代码
CREATE TABLE partition_tbl_2
PARTITION BY (id, dt)
AS SELECT * from partition_tbl_1;

5: 向 Iceberg 表中插入数据

同 StarRocks 内表一致,如果您拥有 Iceberg 表的 INSERT 权限,那么您可以使用 INSERT 将 StarRocks 表数据写入到该 Iceberg 表中(当前仅支持写入到 Parquet 格式的 Iceberg 表)。本功能自 3.1 版本起开始支持。

sql 复制代码
INSERT {INTO | OVERWRITE} <table_name>
[ (column_name [, ...]) ]
{ VALUES ( { expression | DEFAULT } [, ...] ) [, ...] | query }

向指定分区写入数据

sql 复制代码
INSERT {INTO | OVERWRITE} <table_name>
PARTITION (par_col1=<value> [, par_col2=<value>...])
{ VALUES ( { expression | DEFAULT } [, ...] ) [, ...] | query }

向表 partition_tbl_1 中插入如下三行数据

sql 复制代码
INSERT INTO partition_tbl_1
VALUES
    ("buy", 1, "2023-09-01"),
    ("sell", 2, "2023-09-02"),
    ("buy", 3, "2023-09-03");

向表 partition_tbl_1 按指定列顺序插入一个包含简单计算的 SELECT 查询的结果数据

sql 复制代码
INSERT INTO partition_tbl_1 (id, action, dt) SELECT 1+1, 'buy', '2023-09-03';

向表 partition_tbl_1 中插入一个从其自身读取数据的 SELECT 查询的结果数据

sql 复制代码
INSERT INTO partition_tbl_1 SELECT 'buy', 1, date_add(dt, INTERVAL 2 DAY)
FROM partition_tbl_1
WHERE id=1;

向表 partition_tbl_2dt='2023-09-01'id=1 的分区插入一个 SELECT 查询的结果数据

sql 复制代码
INSERT INTO partition_tbl_2 SELECT 'order', 1, '2023-09-01';

sql 复制代码
INSERT INTO partition_tbl_2 partition(dt='2023-09-01',id=1) SELECT 'order';

将表 partition_tbl_1dt='2023-09-01'id=1 的分区下所有 action 列值全部覆盖为 close

sql 复制代码
INSERT OVERWRITE partition_tbl_1 SELECT 'close', 1, '2023-09-01';

sql 复制代码
INSERT OVERWRITE partition_tbl_1 partition(dt='2023-09-01',id=1) SELECT 'close';

6: 基于 External Catalog 的物化视图 - Iceberg

Iceberg 表定义

sql 复制代码
CREATE TABLE spark_catalog.test.iceberg_sample_datetime_day (
  id         BIGINT,
  data       STRING,
  category   STRING,
  ts         TIMESTAMP)
USING iceberg
PARTITIONED BY (days(ts))

基于以上 Iceberg 表创建物化视图

sql 复制代码
CREATE MATERIALIZED VIEW `test_iceberg_datetime_day_mv` (`id`, `data`, `category`, `ts`)
PARTITION BY (`ts`)
DISTRIBUTED BY HASH(`id`)
REFRESH MANUAL
PROPERTIES (
"force_external_table_query_rewrite" = "true"
)
AS
SELECT
  `iceberg_sample_datetime_day`.`id`,
  `iceberg_sample_datetime_day`.`data`,
  `iceberg_sample_datetime_day`.`category`,
  `iceberg_sample_datetime_day`.`ts`
FROM `iceberg`.`test`.`iceberg_sample_datetime_day`;

"force_external_table_query_rewrite" = "true":启用 External Catalog 物化视图的查询改写。StarRocks 支持基于 Hive Catalog、Hudi Catalog、Iceberg Catalog 和 Paimon Catalog 的外部数据源上构建异步物化视图,并支持透明地改写查询。

7: StarRocks 异步物化视图实现增量刷新和透明改写

  • 自动刷新 (REFRESH ASYNC) 的物化视图在基表数据发生变化时会自动刷新。
  • 定时刷新 (REFRESH ASYNC [START (<start_time>)] EVERY (INTERVAL)) 的物化视图将按照定义的间隔定时刷新。
  • 手动刷新 (REFRESH MANUAL) 的物化视图只能通过手动执行 REFRESH MATERIALIZED VIEW 语句进行刷新。您可以指定需要刷新的分区的时间范围,以避免刷新整个物化视图。如果在语句中指定了 FORCE,StarRocks 会强制刷新相应的物化视图或分区,无论基表中的数据是否发生了变化。通过在语句中添加 WITH SYNC MODE,您可以同步调用刷新任务,从而使 StarRocks 仅在任务成功或失败时返回任务结果。
sql 复制代码
CREATE MATERIALIZED VIEW par_mv8
REFRESH ASYNC
PARTITION BY datekey
PROPERTIES(
  "partition_ttl_number" = "2",
  "partition_refresh_number" = "1"
)
AS
SELECT
  k1,
  sum(v1) AS SUM,
  datekey
FROM par_tbl1
GROUP BY datekey, k1;
sql 复制代码
REFRESH MATERIALIZED VIEW par_mv8
PARTITION START ("2021-01-03") END ("2021-01-04")
FORCE WITH SYNC MODE;

大数据相关学习资料

大数据相关学习资料、大数据项目、湖仓一体、架构师必知必会、数据中台建设方法论...

共有1400多份文档资料,另专为星球成员整理了一份比较详细的语雀知识库合计190万字和飞书文档资料。(内容太多,仅展示部分内容...),欢迎大家踊跃加入星球,您将获得:

  1. 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
  2. 免费答疑解惑、交流技术。
  3. 面试指导、模拟面试。
  4. 各类pdf文档下载、星球代码下载。
  5. 提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导。

另外说明加入星球后支持三天无理由退款,不满意无条件随时退。

需要资料请加微信:D1435221412

大数据学习资料

end