Iceberg与StarRocks深度集成实战指南
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 TABLE 或 CREATE 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
包含 id 和 score 两列,如下所示:
sql
CREATE TABLE unpartition_tbl
(
id int,
score double
);
创建分区表 partition_tbl_1
包含 action、id、dt 三列,并定义 id 和 dt 为分区列,如下所示:
sql
CREATE TABLE partition_tbl_1
(
action varchar(20),
id int,
dt date
)
PARTITION BY (id,dt);
查询原表 partition_tbl_1 的数据
并根据查询结果创建分区表 partition_tbl_2,定义 id 和 dt 为 partition_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_2 中 dt='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_1 中 dt='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万字和飞书文档资料。(内容太多,仅展示部分内容...),欢迎大家踊跃加入星球,您将获得:
- 提供最全的大数据知识库,不限设备,随时随地打开看的在线文档。
- 免费答疑解惑、交流技术。
- 面试指导、模拟面试。
- 各类pdf文档下载、星球代码下载。
- 提供简历模板,简历修改指导服务,星球成员免费提供简历修改指导。
另外说明加入星球后支持三天无理由退款,不满意无条件随时退。
需要资料请加微信:D1435221412

end
