Spark SQL 查询引擎连接器

Iceberg

师成师成· 更新于 2026-09-29· 阅读 5 分钟· 0 次阅读

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

Apache Iceberg 是一个面向大规模分析型数据集的开放表格式。Iceberg 通过高性能的表格式为 Spark、Trino、PrestoDB、Flink、Hive 和 Impala 等计算引擎添加表支持,其使用方式与 SQL 表完全相同。

提示

本文假设你已掌握 Iceberg 的基础知识与操作。对于本文中未涉及的 Iceberg 相关知识,你可以查阅其官方文档。

通过使用 Kyuubi,我们可以对 Iceberg 执行 SQL 查询,这比直接使用 Spark 操作 Iceberg 更加便捷、易于理解,也更便于扩展。

Iceberg 集成

要通过 Apache Spark Datasource V2 和 Catalog API 启用 Kyuubi Spark SQL 引擎与 Iceberg 的集成,你需要:

  • 引入 Iceberg 的依赖
  • 设置 Spark 扩展与 Catalog 的配置

依赖

支持 Iceberg 的 Kyuubi Spark SQL 引擎的 classpath 由以下部分组成:

  1. kyuubi-spark-sql-engine-1.9.1_2.12.jar,随 Kyuubi 发行版部署的引擎 jar
  2. 一份 Spark 发行版的副本
  3. iceberg-spark-runtime-<spark.version>_<scala.version>-<iceberg.version>.jar(示例:iceberg-spark-runtime-3.2_2.12-0.14.0.jar),可在 Maven Central 找到

为了让 Iceberg 包对引擎的运行时 classpath 可见,我们可以使用以下方法之一:

  1. 将 Iceberg 包直接放入 $SPARK_HOME/jars
  2. 设置 spark.jars=/path/to/iceberg-spark-runtime

警告

请注意不同 Iceberg 与 Spark 版本之间的兼容性,可在 Iceberg 多引擎支持页面上确认。

配置

要启用 Iceberg 的功能,我们可以设置以下配置:

spark.sql.catalog.spark_catalog=org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.spark_catalog.type=hive
spark.sql.catalog.spark_catalog.uri=thrift://metastore-host:port
spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions

Iceberg 操作

以 CREATE TABLE 为例,

CREATE TABLE foo (
  id bigint COMMENT 'unique id',
  data string)
USING iceberg;

以 SELECT 为例,

SELECT * FROM foo;

以 INSERT 为例,

INSERT INTO foo VALUES (1, 'a'), (2, 'b'), (3, 'c');

以 UPDATE 为例,Spark 3.1 新增了对 UPDATE 查询的支持,用于更新表中匹配的行。

UPDATE foo SET data = 'd', id = 4 WHERE id >= 3 and id < 4;

以 DELETE FROM 为例,Spark 3 新增了对 DELETE FROM 查询的支持,用于从表中删除数据。

DELETE FROM foo WHERE id >= 1 and id < 2;

以 MERGE INTO 为例,

MERGE INTO target_table t
USING source_table s
ON t.id = s.id
WHEN MATCHED AND s.opType = 'delete' THEN DELETE
WHEN MATCHED AND s.opType = 'update' THEN UPDATE SET id = s.id, data = s.data
WHEN NOT MATCHED AND s.opType = 'insert' THEN INSERT (id, data) VALUES (s.id, s.data);

评论

登录后参与评论

正在加载评论…