Apache Spark

快速入门

师成师成· 更新于 2026-09-28· 阅读 14 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

快速开始

Iceberg 的最新版本是 1.11.0。

Spark 目前是功能最完善的 Iceberg 操作计算引擎。我们建议你从 Spark 入手,通过示例来了解 Iceberg 的概念与特性。你也可以在多引擎支持页面查看使用其他计算引擎配合 Iceberg 的文档。

在 Spark 中使用 Iceberg

要在 Spark shell 中使用 Iceberg,请使用 --packages 选项:

spark-shell --packages org.apache.iceberg:iceberg-spark-runtime-4.1_2.13:1.11.0

信息

如果你想将 Iceberg 集成到 Spark 安装中,请将 iceberg-spark-runtime-4.1_2.13 Jar 添加到 Spark 的 jars 文件夹中。

添加目录

Iceberg 自带目录,可让 SQL 命令管理表并按名称加载表。目录通过 spark.sql.catalog.(catalog_name) 下的属性进行配置。

此命令会创建一个名为 local 的基于路径的目录,用于管理 $PWD/warehouse 下的表,并为 Spark 内置目录添加对 Iceberg 表的支持:

spark-sql --packages org.apache.iceberg:iceberg-spark-runtime-4.1_2.13:1.11.0\
    --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions \
    --conf spark.sql.catalog.spark_catalog=org.apache.iceberg.spark.SparkSessionCatalog \
    --conf spark.sql.catalog.spark_catalog.type=hive \
    --conf spark.sql.catalog.local=org.apache.iceberg.spark.SparkCatalog \
    --conf spark.sql.catalog.local.type=hadoop \
    --conf spark.sql.catalog.local.warehouse=$PWD/warehouse

创建表

要在 Spark 中创建你的第一个 Iceberg 表,请使用 spark-sql 命令行或 spark.sql(...) 来执行 CREATE TABLE 命令:

-- local is the path-based catalog defined above
CREATE TABLE local.db.table (id bigint, data string) USING iceberg;
CREATE TABLE source (id bigint, data string) USING parquet;
CREATE TABLE updates (id bigint, data string) USING parquet;

Iceberg 目录支持全套 SQL DDL 命令,包括:

写入

表创建完成后,使用 INSERT INTO 插入数据:

INSERT INTO local.db.table VALUES (1, 'a'), (2, 'b'), (3, 'c');
INSERT INTO source VALUES (10, 'd'), (11, 'ee');
INSERT INTO updates VALUES (1, 'x'), (2, 'x'), (4, 'z');
INSERT INTO local.db.table SELECT id, data FROM source WHERE length(data) = 1;

Iceberg 还为 Spark 增加了行级 SQL 更新功能:MERGE INTO 和 DELETE FROM:

MERGE INTO local.db.table t USING (SELECT * FROM updates) u ON t.id = u.id
WHEN MATCHED THEN UPDATE SET t.data = u.data
WHEN NOT MATCHED THEN INSERT *;

Iceberg 支持使用新的 v2 DataFrame 写入 API 写入 DataFrame:

spark.table("source").select("id", "data")
     .writeTo("local.db.table").append()

旧的 write API 仍受支持,但不推荐使用。

读取

要使用 SQL 进行读取,请在 SELECT 查询中使用 Iceberg 表的名称:

SELECT count(1) as count, data
FROM local.db.table
GROUP BY data;

SQL 也是检查表的推荐方式。要查看表中的所有快照,请使用 snapshots 元数据表:

SELECT * FROM local.db.table.snapshots;
+-------------------------+----------------+-----------+-----------+----------------------------------------------------+-----+
| committed_at            | snapshot_id    | parent_id | operation | manifest_list                                      | ... |
+-------------------------+----------------+-----------+-----------+----------------------------------------------------+-----+
| 2019-02-08 03:29:51.215 | 57897183625154 | null      | append    | s3://.../table/metadata/snap-57897183625154-1.avro | ... |
|                         |                |           |           |                                                    | ... |
|                         |                |           |           |                                                    | ... |
| ...                     | ...            | ...       | ...       | ...                                                | ... |
+-------------------------+----------------+-----------+-----------+----------------------------------------------------+-----+

支持 DataFrame 读取,并且现在可以通过 spark.table 按名称引用表:

val df = spark.table("local.db.table")
df.count()

类型兼容性

Spark 和 Iceberg 支持不同的类型集合。Iceberg 会自动进行类型转换,但并非适用于所有组合,因此在设计表中各列的类型之前,你可能需要先了解 Iceberg 中的类型转换规则。

Spark 类型到 Iceberg 类型

该类型转换表描述了 Spark 类型如何转换为 Iceberg 类型。此转换同时适用于创建 Iceberg 表以及通过 Spark 向 Iceberg 表写入数据的场景。

SparkIceberg备注
booleanboolean
shortinteger
byteinteger
integerinteger
longlong
floatfloat
doubledouble
datedate
timestamptimestamp with timezone
timestamp_ntztimestamp without timezone
charstring
varcharstring
stringstring
binarybinary
decimaldecimal
structstruct
arraylist
mapmap

说明

该表是基于建表时的转换表示的。实际上,写入时支持更广泛的转换。以下是关于写入的几点说明:

  • Iceberg 的数值类型(integer、long、float、double、decimal)在写入时支持类型提升。例如,你可以将 Spark 类型 short、byte、integer、long 写入 Iceberg 类型 long。
  • 你可以使用 Spark 的 binary 类型写入 Iceberg 的 fixed 类型。注意,写入时会进行长度校验。

Iceberg 类型到 Spark 类型

该类型转换表描述了 Iceberg 类型如何转换为 Spark 类型。此转换适用于通过 Spark 从 Iceberg 表读取数据的场景。

IcebergSpark备注
booleanboolean
integerinteger
longlong
floatfloat
doubledouble
datedate
time不支持
timestamp with timezonetimestamp
timestamp without timezonetimestamp_ntz
stringstring
uuidstring
fixedbinary
binarybinary
decimaldecimal
structstruct
listarray
mapmap
nanosecond timestamp不支持
nanosecond timestamp with timezone不支持
unknownnullSpark 4.0+
variantvariantSpark 4.0+
geometry不支持
geography不支持

下一步

接下来,你可以进一步了解 Spark 中的 Iceberg 表:

评论

登录后参与评论

正在加载评论…