改版通知

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

Iceberg与Spark深度集成实战指南

ckckck2025年1月10日7 浏览

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 COLUMNDROP 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