Iceberg与Spark深度集成实战指南
1 环境准备
1.1 安装 Spark
Spark 使用的是 3.3.1 版本,Iceberg 使用的是 1.1.0 版本。
1) Spark 与 Iceberg 的版本对应关系
| Spark 版本 | Iceberg 版本 |
|---|---|
| 2.4 | 0.7.0-incubating – 1.1.0 |
| 3 | 0.9.0 – 1.0.0 |
| 3.1 | 0.12.0 – 1.1.0 |
| 3.2 | 0.13.0 – 1.1.0 |
| 3.3 | 0.14.0 – 1.1.0 |
0.12.1 支持 Spark 2.4+,但不完善;建议使用 Spark 3.x 以上版本。
2) 上传并解压 Spark 安装包
bash
tar -zxvf spark-3.3.1-bin-hadoop3.tgz
mv spark-3.3.1-bin-hadoop3 spark-3.3.1
3) 配置环境变量
bash
sudo vim /etc/profile.d/my_env.sh
export SPARK_HOME=/opt/module/spark-3.3.1
export PATH=$PATH:$SPARK_HOME/bin
source /etc/profile.d/my_env.sh
4) 拷贝 Iceberg 的 jar 包到 Spark 的 jars 目录
bash
cp /opt/software/iceberg/iceberg-spark-runtime-3.3_2.12-1.1.0.jar /opt/module/spark-3.3.1/jars
1.2 启动 Hadoop
2 Spark 配置 Catalog
Spark 中支持两种 Catalog 的设置:Hive 和 Hadoop。Hive Catalog 是 Iceberg 表存储使用 Hive 默认的数据路径,Hadoop Catalog 需要指定 Iceberg 格式表存储路径。
bash
vim spark-defaults.conf
2.1 Hive Catalog
bash
spark.sql.catalog.hive_prod = org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.hive_prod.type = hive
spark.sql.catalog.hive_prod.uri = thrift://192.168.110.120:9083
use hive_prod.db;
2.2 Hadoop Catalog
bash
spark.sql.catalog.hadoop_prod = org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.hadoop_prod.type = hadoop
spark.sql.catalog.hadoop_prod.warehouse = hdfs://192.168.110.120:8020/warehouse/spark-iceberg
use hadoop_prod.db;
3 启动 Spark SQL
bash
./bin/spark-sql
4 SQL 操作
4.1 创建表
使用 Catalog hadoop_prod
sql
spark-sql> use hadoop_prod;
-- 创建数据库,默认没有
spark-sql> create database default;
spark-sql> use default;
-- 使用 hive_prod catalog,默认库是 default
spark-sql> use hive_prod;
spark-sql> use default;
spark-sql> CREATE TABLE hadoop_prod.default.sample1 (
id bigint COMMENT 'unique id',
data string
) USING iceberg;
hadoop_prod: hadoop_catalog 名称USING iceberg: 创建为 Iceberg 表PARTITIONED BY (partition-expressions): 配置分区LOCATION '(fully-qualified-uri)': 指定表路径COMMENT 'table documentation': 配置表备注TBLPROPERTIES ('key'='value', ...): 配置表属性
表属性可以参考 Iceberg 官网文档。
对 Iceberg 表的每次更改都会生成一个新的元数据文件(json 文件)以提供原子性。默认情况下,旧元数据文件作为历史文件保存不会删除。如果要自动清除元数据文件,在表属性中设置 write.metadata.delete-after-commit.enabled=true。这将保留一些元数据文件(直到 write.metadata.previous-versions-max),并在每个新创建的元数据文件之后删除旧的元数据文件。
1) 创建分区表
1) 分区表
sql
CREATE TABLE hadoop_prod.default.sample2 (
id bigint,
data string,
category string
) USING iceberg
PARTITIONED BY (category);
2) 创建隐藏分区表
sql
CREATE TABLE hadoop_prod.default.sample3 (
id bigint,
data string,
category string,
ts timestamp
) USING iceberg
PARTITIONED BY (bucket(16, id), days(ts), category);
支持的转换有:
years(ts): 按年划分months(ts): 按月划分days(ts)或date(ts): 等效于 dateint 分区hours(ts)或date_hour(ts): 等效于 dateint 和 hour 分区bucket(N, col): 按哈希值划分 mod N 个桶truncate(L, col): 按截断为 L 的值划分
隐藏分区:可以指定分区字段做计算,计算的结果不用体现在字段定义中。
2) 使用 CTAS 语法建表
sql
CREATE TABLE hadoop_prod.default.sample4
USING iceberg
AS SELECT * FROM hadoop_prod.default.sample3;
不指定分区就是无分区,需要重新指定分区、表属性:
sql
CREATE TABLE hadoop_prod.default.sample5
USING iceberg
PARTITIONED BY (bucket(8, id), hours(ts), category)
TBLPROPERTIES ('key'='value')
AS SELECT * FROM hadoop_prod.default.sample3;
3) 使用 Replace table 建表
sql
REPLACE TABLE hadoop_prod.default.sample5
USING iceberg
AS SELECT * FROM hadoop_prod.default.sample3;
REPLACE TABLE hadoop_prod.default.sample5
USING iceberg
PARTITIONED BY (part)
TBLPROPERTIES ('key'='value')
AS SELECT * FROM hadoop_prod.default.sample3;
-- 存在替换,不存在创建
CREATE OR REPLACE TABLE hadoop_prod.default.sample6
USING iceberg
AS SELECT * FROM hadoop_prod.default.sample3;
4.2 删除表
对于 HadoopCatalog 而言,运行 DROP TABLE 将从 catalog 中删除表并删除表内容。
sql
CREATE EXTERNAL TABLE hadoop_prod.default.sample7 (
id bigint COMMENT 'unique id',
data string
) USING iceberg;
INSERT INTO hadoop_prod.default.sample7 VALUES (1, 'a');
DROP TABLE hadoop_prod.default.sample7;
对于 HiveCatalog 而言:
- 在 0.14 之前,运行
DROP TABLE将从 catalog 中删除表并删除表内容。 - 从 0.14 开始,
DROP TABLE只会从 catalog 中删除表,不会删除数据。为了删除表内容,应该使用DROP TABLE PURGE。
sql
CREATE TABLE hive_prod.default.sample7 (
id bigint COMMENT 'unique id',
data string
) USING iceberg;
INSERT INTO hive_prod.default.sample7 VALUES (1, 'a');
1) 删除表
sql
DROP TABLE hive_prod.default.sample7;
2) 删除表和数据
sql
DROP TABLE hive_prod.default.sample7 PURGE;
4.3 修改表
Iceberg 在 Spark 3(只在 Spark 3 中支持)中完全支持 ALTER TABLE,包括:
- 重命名表
- 设置或删除表属性
- 添加、删除和重命名列
- 添加、删除和重命名嵌套字段
- 重新排序顶级列和嵌套结构字段
- 扩大 int、float 和 decimal 字段的类型
- 将必选列变为可选列
此外,还可以使用 SQL 扩展来添加对分区演变的支持和设置表的写顺序。
sql
CREATE TABLE hive_prod.default.sample1 (
id bigint COMMENT 'unique id',
data string
) USING iceberg;
1) 修改表名(不支持修改 HadoopCatalog 的表名)
sql
ALTER TABLE hive_prod.default.sample1 RENAME TO hive_prod.default.sample2;
2) 修改表属性
(1) 修改表属性
sql
ALTER TABLE hive_prod.default.sample1 SET TBLPROPERTIES ('read.split.target-size'='268435456');
ALTER TABLE hive_prod.default.sample1 SET TBLPROPERTIES ('comment' = 'A table comment.');
(2) 删除表属性
sql
ALTER TABLE hive_prod.default.sample1 UNSET TBLPROPERTIES ('read.split.target-size');
3) 添加列
sql
ALTER TABLE hadoop_prod.default.sample1
ADD COLUMNS (category string COMMENT 'new_column');
-- 添加 struct 类型的列
ALTER TABLE hadoop_prod.default.sample1
ADD COLUMN point struct<x: double, y: double>;
-- 往 struct 类型的列中添加字段
ALTER TABLE hadoop_prod.default.sample1
ADD COLUMN point.z double;
-- 创建 struct 的嵌套数组列
ALTER TABLE hadoop_prod.default.sample1
ADD COLUMN points array<struct<x: double, y: double>>;
-- 在数组中的结构中添加一个字段。使用关键字 'element' 访问数组的元素列。
ALTER TABLE hadoop_prod.default.sample1
ADD COLUMN points.element.z double;
-- 创建一个包含 Map 类型的列,key 和 value 都为 struct 类型
ALTER TABLE hadoop_prod.default.sample1
ADD COLUMN pointsm map<struct<x: int>, struct<a: int>>;
-- 在 Map 类型的 value 的 struct 中添加一个字段。
ALTER TABLE hadoop_prod.default.sample1
ADD COLUMN pointsm.value.b int;
-- 在 Spark 2.4.4 及以后版本中,可以通过添加 FIRST 或 AFTER 子句在任何位置添加列
ALTER TABLE hadoop_prod.default.sample1
ADD COLUMN new_column1 bigint AFTER id;
ALTER TABLE hadoop_prod.default.sample1
ADD COLUMN new_column2 bigint FIRST;
4) 修改列
(1) 修改列名
sql
ALTER TABLE hadoop_prod.default.sample1 RENAME COLUMN data TO data1;
(2) Alter Column 修改类型(只允许安全的转换)
sql
ALTER TABLE hadoop_prod.default.sample1
ADD COLUMNS (idd int);
ALTER TABLE hadoop_prod.default.sample1 ALTER COLUMN idd TYPE bigint;
(3) Alter Column 修改列的注释
sql
ALTER TABLE hadoop_prod.default.sample1 ALTER COLUMN id TYPE double COMMENT 'a';
ALTER TABLE hadoop_prod.default.sample1 ALTER COLUMN id COMMENT 'b';
(4) Alter Column 修改列的顺序
sql
ALTER TABLE hadoop_prod.default.sample1 ALTER COLUMN id FIRST;
ALTER TABLE hadoop_prod.default.sample1 ALTER COLUMN new_column2 AFTER new_column1;
(5) Alter Column 修改列是否允许为 null
sql
ALTER TABLE hadoop_prod.default.sample1 ALTER COLUMN id DROP NOT NULL;
ALTER COLUMN 不用于更新 struct 类型。使用 ADD COLUMN 和 DROP COLUMN 添加或删除 struct 类型的字段。
5) 删除列
sql
ALTER TABLE hadoop_prod.default.sample1 DROP COLUMN idd;
ALTER TABLE hadoop_prod.default.sample1 DROP COLUMN point.z;
6) 添加分区(Spark 2.4 不支持,Spark 3 需要配置扩展)
bash
vim spark-default.conf
spark.sql.extensions = org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions
重新进入 spark-sql shell:
sql
ALTER TABLE hadoop_prod.default.sample1 ADD PARTITION FIELD category;
ALTER TABLE hadoop_prod.default.sample1 ADD PARTITION FIELD bucket(16, id);
ALTER TABLE hadoop_prod.default.sample1 ADD PARTITION FIELD truncate(data, 4);
ALTER TABLE hadoop_prod.default.sample1 ADD PARTITION FIELD years(ts);
ALTER TABLE hadoop_prod.default.sample1 ADD PARTITION FIELD bucket(16, id) AS shard;
7) 删除分区(Spark 3 需要配置扩展)
sql
ALTER TABLE hadoop_prod.default.sample1 DROP PARTITION FIELD category;
ALTER TABLE hadoop_prod.default.sample1 DROP PARTITION FIELD bucket(16, id);
ALTER TABLE hadoop_prod.default.sample1 DROP PARTITION FIELD truncate(data, 4);
ALTER TABLE hadoop_prod.default.sample1 DROP PARTITION FIELD years(ts);
ALTER TABLE hadoop_prod.default.sample1 DROP PARTITION FIELD shard;
注意:尽管删除了分区,但列仍然存在于表结构中。删除分区字段是元数据操作,不会改变任何现有的表数据。新数据将被写入新的分区,但现有数据将保留在旧的分区布局中。当分区发生变化时,动态分区覆盖行为也会发生变化。例如,如果按天划分分区,而改为按小时划分分区,那么覆盖将覆盖每小时划分的分区,而不再覆盖按天划分的分区。删除分区字段时要小心,可能导致元数据查询失败或产生不同的结果。
8) 修改分区(Spark 3 需要配置扩展)
sql
ALTER TABLE hadoop_prod.default.sample1 REPLACE PARTITION FIELD bucket(16, id) WITH bucket(8, id);
9) 修改表的写入顺序
sql
ALTER TABLE hadoop_prod.default.sample1 WRITE ORDERED BY category, id;
ALTER TABLE hadoop_prod.default.sample1 WRITE ORDERED BY category ASC, id DESC;
ALTER TABLE hadoop_prod.default.sample1 WRITE ORDERED BY category ASC NULLS LAST, id DESC NULLS FIRST;
表写顺序不能保证查询的数据顺序。它只影响数据写入表的方式。WRITE ORDERED BY 设置了一个全局排序,即跨任务的行排序,就像在 INSERT 命令中使用 ORDER BY 一样:
sql
INSERT INTO hadoop_prod.default.sample1
SELECT id, data, category, ts FROM another_table
ORDER BY ts, category;
要在每个任务内排序,而不是跨任务排序,使用 local ORDERED BY:
sql
ALTER TABLE hadoop_prod.default.sample1 WRITE LOCALLY ORDERED BY category, id;
10) 按分区并行写入
sql
ALTER TABLE hadoop_prod.default.sample1 WRITE DISTRIBUTED BY PARTITION;
ALTER TABLE hadoop_prod.default.sample1 WRITE DISTRIBUTED BY PARTITION LOCALLY ORDERED BY category, id;
4.4 插入数据
sql
CREATE TABLE hadoop_prod.default.a (id bigint, count bigint)
USING iceberg;
CREATE TABLE hadoop_prod.default.b (id bigint, count bigint, flag string)
USING iceberg;
1) Insert Into
sql
INSERT INTO hadoop_prod.default.a VALUES (1, 1), (2, 2), (3, 3);
INSERT INTO hadoop_prod.default.b VALUES (1, 1, 'a'), (2, 2, 'b'), (4, 4, 'd');
2) MERGE INTO 行级更新
sql
MERGE INTO hadoop_prod.default.a t
USING (SELECT * FROM hadoop_prod.default.b) u ON t.id = u.id
WHEN MATCHED AND u.flag='b' THEN UPDATE SET t.count = t.count + u.count
WHEN MATCHED AND u.flag='a' THEN DELETE
WHEN NOT MATCHED THEN INSERT (id, count) VALUES (u.id, u.count);
3) INSERT OVERWRITE
INSERT OVERWRITE 可以覆盖 Iceberg 表中的数据,这种操作会将表中全部数据替换掉,建议如果有部分数据替换操作可以使用 MERGE INTO 操作。对于 Iceberg 分区表使用 INSERT OVERWRITE 操作时,有两种情况,第一种是“动态覆盖”,第二种是“静态覆盖”。
- 动态分区覆盖:动态覆盖会全量将原有数据覆盖,并将新插入的数据根据 Iceberg 表分区规则自动分区,类似 Hive 中的动态分区。
- 静态分区覆盖:静态覆盖需要在向 Iceberg 中插入数据时需要手动指定分区,如果当前 Iceberg 表存在这个分区,那么只有这个分区的数据会被覆盖,其他分区数据不受影响,如果 Iceberg 表不存在这个分区,那么相当于给 Iceberg 表增加了一个分区。
sql
-- 动态分区覆盖
INSERT OVERWRITE hadoop_prod.default.test1
SELECT id, name, loc FROM hadoop_prod.default.test3;
-- 静态分区覆盖
INSERT OVERWRITE hadoop_prod.default.test1
PARTITION (loc = "jiangsu")
SELECT id, name FROM hadoop_prod.default.test3;
4) DELETE FROM
Spark 3.x 版本之后支持 DELETE FROM 可以根据指定的 where 条件来删除表中数据。如果 where 条件匹配 Iceberg 表一个分区的数据,Iceberg 仅会修改元数据,如果 where 条件匹配的表的单个行,则 Iceberg 会重写受影响行所在的数据文件。
sql
-- 创建表 delete_tbl,并加载数据
spark.sql(
"""
|CREATE TABLE hadoop_prod.default.delete_tbl (id int, name string, age int) USING iceberg
""".stripMargin)
spark.sql(
"""
|INSERT INTO hadoop_prod.default.delete_tbl VALUES (1, "zs", 18), (2, "ls", 19), (3, "ww", 20), (4, "ml", 21), (5, "tq", 22), (6, "gb", 23)
""".stripMargin)
-- 根据条件范围删除表 delete_tbl 中的数据
spark.sql(
"""
|DELETE FROM hadoop_prod.default.delete_tbl WHERE id > 3 AND id < 6
""".stripMargin)
spark.sql("SELECT * FROM hadoop_prod.default.delete_tbl").show()
-- 根据等值条件删除表 delete_tbl 中的一条数据
spark.sql(
"""
|DELETE FROM hadoop_prod.default.delete_tbl WHERE id = 2
""".stripMargin)
5) UPDATE
Spark 3.x+ 版本支持了 UPDATE 更新数据操作,可以根据匹配的条件进行数据更新操作。
sql
UPDATE hadoop_prod.default.update_tbl SET name = 'zhangsan', age = 30 WHERE id <= 3;
4.5 查询数据
1) 普通查询
sql
SELECT COUNT(1) AS count, data
FROM local.db.table
GROUP BY data;
2) 查询元数据
sql
-- 查询表快照
SELECT * FROM hadoop_prod.default.a.snapshots;
-- 查询数据文件信息
SELECT * FROM hadoop_prod.default.a.files;
-- 查询表历史
SELECT * FROM hadoop_prod.default.a.history;
-- 查询 manifest
SELECT * FROM hadoop_prod.default.a.manifests;
4.6 存储过程
Procedures 可以通过 CALL 从任何已配置的 Iceberg Catalog 中使用。所有 Procedures 都在 namespace 中。
1) 参数传递
按照参数名传参:
sql
CALL catalog_name.system.procedure_name(arg_name_2 => arg_2, arg_name_1 => arg_1);
当按位置传递参数时,如果结束参数是可选的,则只有结束参数可以省略。
sql
CALL catalog_name.system.procedure_name(arg_1, arg_2, ... arg_n);
2) 快照管理
(1) 回滚到指定的快照 id
sql
CALL hadoop_prod.system.rollback_to_snapshot('default.a', 7601163594701794741);
(2) 回滚到指定时间的快照
sql
CALL hadoop_prod.system.rollback_to_timestamp('db.sample', TIMESTAMP '2021-06-30 00:00:00.000');
(3) 设置表的当前快照 ID
sql
CALL hadoop_prod.system.set_current_snapshot('db.sample', 1);
(4) 从快照变为当前表状态
sql
CALL hadoop_prod.system.cherrypick_snapshot('default.a', 7629160535368763452);
CALL hadoop_prod.system.cherrypick_snapshot(snapshot_id => 7629160535368763452, table => 'default.a');
3) 元数据管理
(1) 删除早于指定日期和时间的快照,但保留最近 100 个快照
sql
CALL hive_prod.system.expire_snapshots('db.sample', TIMESTAMP '2021-06-30 00:00:00.000', 100);
(2) 删除 Iceberg 表中任何元数据文件中没有引用的文件
列出所有需要删除的候选文件
sql
CALL catalog_name.system.remove_orphan_files(table => 'db.sample', dry_run => true);
删除指定目录中 db.sample 表不知道的任何文件
sql
CALL catalog_name.system.remove_orphan_files(table => 'db.sample', location => 'tablelocation/data');
(3) 合并数据文件(合并小文件)
sql
CALL catalog_name.system.rewrite_data_files('db.sample');
CALL catalog_name.system.rewrite_data_files(table => 'db.sample', strategy => 'sort', sort_order => 'id DESC NULLS LAST, name ASC NULLS FIRST');
CALL catalog_name.system.rewrite_data_files(table => 'db.sample', strategy => 'sort', sort_order => 'zorder(c1,c2)');
CALL catalog_name.system.rewrite_data_files(table => 'db.sample', options => map('min-input-files','2'));
CALL catalog_name.system.rewrite_data_files(table => 'db.sample', where => 'id = 3 and name = "foo"');
(4) 重写表清单来优化执行计划
sql
CALL catalog_name.system.rewrite_manifests('db.sample');
重写表 db 中的清单,并禁用 Spark 缓存的使用。这样做可以避免执行程序上的内存问题。
sql
CALL catalog_name.system.rewrite_manifests('db.sample', false);
4) 迁移表
(1) 快照
sql
CALL catalog_name.system.snapshot('db.sample', 'db.snap');
CALL catalog_name.system.snapshot('db.sample', 'db.snap', '/tmp/temptable/');
(2) 迁移
sql
CALL catalog_name.system.migrate('spark_catalog.db.sample', map('foo', 'bar'));
CALL catalog_name.system.migrate('db.sample');
(3) 添加数据文件
sql
CALL spark_catalog.system.add_files(
table => 'db.tbl',
source_table => 'db.src_tbl',
partition_filter => map('part_col_1', 'A')
);
CALL spark_catalog.system.add_files(table => 'db.tbl',
source_table => 'parquet.path/to/table'
);
5) 元数据信息
(1) 获取指定快照的父快照 id
sql
CALL spark_catalog.system.ancestors_of('db.tbl');
(2) 获取指定快照的所有祖先快照
sql
CALL spark_catalog.system.ancestors_of('db.tbl', 1);
CALL spark_catalog.system.ancestors_of(snapshot_id => 1, table => 'db.tbl');
5 DataFrame 操作
5.1 环境准备
1) 创建 Maven 工程,配置 pom 文件
xml
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.atguigu.iceberg</groupId>
<artifactId>spark-iceberg-demo</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<scala.binary.version>2.12</scala.binary.version>
<spark.version>3.3.1</spark.version>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
</properties>
<dependencies>
<!-- Spark 的依赖引入 -->
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-core_${scala.binary.version}</artifactId>
<version>${spark.version}</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-sql_${scala.binary.version}</artifactId>
<version>${spark.version}</version>
</dependency>
<dependency>
<groupId>org.apache.spark</groupId>
<artifactId>spark-hive_${scala.binary.version}</artifactId>
<version>${spark.version}</version>
</dependency>
<!-- fastjson <= 1.2.80 存在安全漏洞 -->
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>1.2.83</version>
</dependency>
<!-- Iceberg 的依赖引入 -->
<dependency>
<groupId>org.apache.iceberg</groupId>
<artifactId>iceberg-spark-runtime-3.3_2.12</artifactId>
<version>1.1.0</version>
</dependency>
</dependencies>
<build>
<plugins>
<!-- assembly 打包插件 -->
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-assembly-plugin</artifactId>
<version>3.0.0</version>
<executions>
<execution>
<id>make-assembly</id>
<phase>package</phase>
<goals>
<goal>single</goal>
</goals>
</execution>
</executions>
<configuration>
<archive>
<manifest>
</manifest>
</archive>
<descriptorRefs>
<descriptorRef>jar-with-dependencies</descriptorRef>
</descriptorRefs>
</configuration>
</plugin>
<!-- Maven 编译 scala 所需依赖 -->
<plugin>
<groupId>net.alchim31.maven</groupId>
<artifactId>scala-maven-plugin</artifactId>
<version>3.2.2</version>
<executions>
<execution>
<goals>
<goal>compile</goal>
<goal>testCompile</goal>
</goals>
</execution>
</executions>
</plugin>
</plugins>
</build>
</project>
2) 配置 Catalog
scala
val spark: SparkSession = SparkSession.builder()
.master("local")
.appName(this.getClass.getSimpleName)
// 指定 hive catalog,catalog 名称为 iceberg_hive
.config("spark.sql.catalog.iceberg_hive", "org.apache.iceberg.spark.SparkCatalog")
.config("spark.sql.catalog.iceberg_hive.type", "hive")
.config("spark.sql.catalog.iceberg_hive.uri", "thrift://hadoop1:9083")
// 指定 hadoop catalog,catalog 名称为 iceberg_hadoop
.config("spark.sql.catalog.iceberg_hadoop", "org.apache.iceberg.spark.SparkCatalog")
.config("spark.sql.catalog.iceberg_hadoop.type", "hadoop")
.config("spark.sql.catalog.iceberg_hadoop.warehouse", "hdfs://hadoop1:8020/warehouse/spark-iceberg")
.getOrCreate()
5.2 读取表
1) 加载表
scala
spark.read
.format("iceberg")
.load("hdfs://hadoop1:8020/warehouse/spark-iceberg/default/a")
.show()
或
scala
// 仅支持 Spark 3.0 以上
spark.table("iceberg_hadoop.default.a").show()
或
scala
val spark: SparkSession = SparkSession.builder()
.master("local")
.appName("test")
// 指定 hadoop catalog,catalog 名称为 hadoop_prod
.config("spark.sql.catalog.hadoop_prod", "org.apache.iceberg.spark.SparkCatalog")
.config("spark.sql.catalog.hadoop_prod.type", "hadoop")
.config("spark.sql.catalog.hadoop_prod.warehouse", "hdfs://mycluster/sparkoperateiceberg")
.config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
.getOrCreate()
// 1. 创建 Iceberg 表,并插入数据
spark.sql(
"""
|CREATE TABLE hadoop_prod.mydb.mytest (id int, name string, age int) USING iceberg
""".stripMargin)
spark.sql(
"""
|INSERT INTO hadoop_prod.mydb.mytest VALUES (1, "zs", 18), (2, "ls", 19), (3, "ww", 20)
""".stripMargin)
// 1. SQL 方式读取 Iceberg 中的数据
spark.sql("SELECT * FROM hadoop_prod.mydb.mytest").show()
// 2. 使用 Spark 查询 Iceberg 中的表除了使用 SQL 方式之外,还可以使用 DataFrame 方式,建议使用 SQL 方式
// 第一种方式使用 DataFrame 方式查询 Iceberg 表数据
val frame1: DataFrame = spark.table("hadoop_prod.mydb.mytest")
frame1.show()
// 第二种方式使用 DataFrame 加载 Iceberg 表数据
val frame2: DataFrame = spark.read.format("iceberg").load("hdfs://mycluster/sparkoperateiceberg/mydb/mytest")
frame2.show()
2) 时间旅行:指定时间查询
scala
spark.read
.option("as-of-timestamp", "499162860000")
.format("iceberg")
.load("hdfs://192.168.110.120:8020/warehouse/spark-iceberg/default/a")
.show()
3) 时间旅行:指定快照 id 查询
scala
spark.read
.option("snapshot-id", "7601163594701794741")
.format("iceberg")
.load("hdfs://192.168.110.120:8020/warehouse/spark-iceberg/default/a")
.show()
4) 增量查询
scala
spark.read
.format("iceberg")
.option("start-snapshot-id", "10963874102873")
.option("end-snapshot-id", "63874143573109")
.load("hdfs://192.168.110.120:8020/warehouse/spark-iceberg/default/a")
.show()
查询的表只能是 append 的方式写数据,不支持 replace, overwrite, delete 操作。
5) 其他常见操作命令
sql
-- 查询快照
SELECT * FROM hadoop_prod.mydb.mytest.snapshots;
-- 查询表历史
SELECT * FROM hadoop_prod.mydb.mytest.history;
-- 查询表 data files
SELECT * FROM hadoop_prod.mydb.mytest.files;
-- 查询 Manifests
SELECT * FROM hadoop_prod.mydb.mytest.manifests;
-- 查询指定快照数据
spark.read
.option("snapshot-id", 3368002881426159310L)
.format("iceberg")
.load("hdfs://mycluster/sparkoperateiceberg/mydb/mytest")
.show();
-- 根据时间戳查找
spark.read.option("as-of-timestamp", "1640066148000")
.format("iceberg")
.load("hdfs://mycluster/sparkoperateiceberg/mydb/mytest")
.show();
-- 根据时间戳查找
CALL hadoop_prod.system.rollback_to_timestamp('mydb.mytest', TIMESTAMP '2021-12-23 16:56:40.000');
-- 回滚快照
val conf = new Configuration()
val catalog = new HadoopCatalog(conf, "hdfs://mycluster/sparkoperateiceberg")
catalog.setConf(conf)
val table: Table = catalog.loadTable(TableIdentifier.of("mydb", "mytest"))
table.manageSnapshots().rollbackTo(3368002881426159310L).commit();
-- 回滚快照
CALL hadoop_prod.system.rollback_to_snapshot("mydb.mytest", 5440886662709904549);
-- 删除历史快照
table.expireSnapshots().expireOlderThan(1640070000000L).commit();
-- 删除早于某个时间的快照,但保留最近 N 个快照
CALL ${Catalog 名称}.system.expire_snapshots("${库名.表名}", TIMESTAMP '年-月-日 时-分-秒.000', N);
5.3 写入表
1) 插入数据并建表
scala
package org.scala.iceberg.spark
import org.apache.spark.sql.{DataFrame, SparkSession}
object Demo1 {
def main(args: Array[String]): Unit = {
// 写入表报错:Permssion denied: user=Administrator, access=WRITE, inode="/warehouse/spark-iceberg/default"
// 可以在最上方加入下边这行代码设置 HADOOP_PROXY_USER 环境变量为 root 用户
System.setProperty("HADOOP_USER_NAME", "root")
// 1. 创建 catalog
val spark: SparkSession = SparkSession.builder()
.master("local[*]")
.appName(this.getClass.getSimpleName)
// 设置 hive_catalog,catalog_name 为 iceberg_hive
.config("spark.sql.catalog.iceberg_hive", "org.apache.iceberg.spark.SparkCatalog")
.config("spark.sql.catalog.iceberg_hive.type", "hive")
.config("spark.sql.catalog.iceberg_hive.uri", "thrift://192.168.110.120:9083")
// 支持 iceberg
.config("iceberg.engine.hive.enabled", "true")
// 设置 hive_catalog,catalog_name 为 iceberg_hive
.config("spark.sql.catalog.iceberg_hadoop", "org.apache.iceberg.spark.SparkCatalog")
.config("spark.sql.catalog.iceberg_hadoop
end
