Spark SQL 查询引擎连接器
Iceberg
登录后可跨设备保存划线和私人笔记登录
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 的 Kyuubi Spark SQL 引擎的 classpath 由以下部分组成:
- kyuubi-spark-sql-engine-1.9.1_2.12.jar,随 Kyuubi 发行版部署的引擎 jar
- 一份 Spark 发行版的副本
- 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 可见,我们可以使用以下方法之一:
- 将 Iceberg 包直接放入
$SPARK_HOME/jars - 设置
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.IcebergSparkSessionExtensionsIceberg 操作
以 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);评论
登录后参与评论
正在加载评论…
KnowForge