Nessie
Iceberg Nessie 集成
Iceberg 通过 iceberg-nessie 模块提供与 Nessie 的集成。本节介绍如何在 Iceberg 中使用 Nessie。Nessie 在 Iceberg 之上提供了若干关键特性:
- 跨表事务
- 类似 git 的操作(例如分支、标签、提交)
- 类似 Hive 的元数据存储能力
有关 Nessie 的更多信息,请参阅 Project Nessie。Nessie 需要运行一个服务端,请参阅 Getting Started 来启动 Nessie 服务。
启用 Nessie Catalog
从 0.11.0 开始的所有版本中,iceberg-nessie 模块都已随 Spark 和 Flink 运行时一起打包。要开始使用 Nessie(配合 spark-3.5)和 Iceberg,只需将 Iceberg 运行时添加到你的进程中即可。例如:spark-sql --packages org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.11.0。
Spark SQL 扩展
可以使用 Nessie SQL 扩展来管理 Nessie 仓库,如下所示。以 Spark 3.5 和 scala 2.12 为例:
bin/spark-sql
--packages "org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.11.0,org.projectnessie.nessie-integrations:nessie-spark-extensions-3.5_2.12:0.105.3"
--conf spark.sql.extensions="org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions,org.projectnessie.spark.extensions.NessieSparkSessionExtensions"
--conf <other settings>请参阅 Nessie SQL 扩展文档 以了解更多信息。
Nessie Catalog
0.11.0 版本引入的一项主要功能,是能够从 Spark 和 Flink 轻松访问自定义 Catalog。有关向 Iceberg 添加自定义 Catalog 的说明,请参阅 Spark 配置和 Flink 配置。
要使用 Nessie Catalog,需要以下属性:
warehouse。与大多数其他 Catalog 一样,warehouse 属性是一个文件路径,用于指定该 Catalog 存储表的位置。uri。这是 Nessie 服务器的基础 uri。例如http://localhost:19120/api/v2。ref(可选)。这是你要操作的 Nessie 分支或标签。
直接在 Java 中运行时,代码如下:
Map<String, String> options = new HashMap<>();
options.put("warehouse", "/path/to/warehouse");
options.put("ref", "main");
options.put("uri", "https://localhost:19120/api/v2");
Catalog nessieCatalog = CatalogUtil.loadCatalog("org.apache.iceberg.nessie.NessieCatalog", "nessie", options, hadoopConfig);以及在 Spark 中:
conf.set("spark.sql.catalog.nessie.warehouse", "/path/to/warehouse");
conf.set("spark.sql.catalog.nessie.uri", "http://localhost:19120/api/v2")
conf.set("spark.sql.catalog.nessie.ref", "main")
conf.set("spark.sql.catalog.nessie.type", "nessie")
conf.set("spark.sql.catalog.nessie", "org.apache.iceberg.spark.SparkCatalog")
conf.set("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions,org.projectnessie.spark.extensions.NessieSparkSessionExtensions")以下是通过 Python API 在 Flink 中的使用方式(更多详细信息请参见此处):
import os
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment
env = StreamExecutionEnvironment.get_execution_environment()
iceberg_flink_runtime_jar = os.path.join(os.getcwd(), "iceberg-flink-runtime-1.11.0.jar")
env.add_jars("file://{}".format(iceberg_flink_runtime_jar))
table_env = StreamTableEnvironment.create(env)
table_env.execute_sql("CREATE CATALOG nessie_catalog WITH ("
"'type'='iceberg', "
"'type'='nessie', "
"'uri'='http://localhost:19120/api/v2', "
"'ref'='main', "
"'warehouse'='/path/to/warehouse')")上面提到的 nessie 这个名称本身并没有什么特殊之处。Spark catalog 可以使用任意名称,关键在于 type 或 catalog-impl 的设置,以及正确启动 Nessie 所需的配置。一旦拥有了 Nessie catalog,你就可以访问整个 Nessie 仓库。随后,你可以对分支执行创建/删除/合并操作,并在分支上提交。
Nessie Catalog 中的每个 Iceberg 表都由一个任意长度的命名空间和表名来标识(例如 data.base.name.table)。如此处所述,这些命名空间必须被显式创建。对启用 Nessie 的 Iceberg 表执行的任何事务,在 Nessie 中都是一次提交。Nessie 的一次提交可以包含对任意数量的表执行的任意数量的操作,但在 Iceberg 中,这将受限于当前可用的单表事务集合。
其他操作(如合并、查看提交日志或 diff)则通过在 Java 中直接与 NessieClient 交互,或使用 Python 客户端或命令行工具来完成。有关命令行工具的更多详情,请参阅 Nessie CLI;有关 Nessie 功能的更完整说明,请参阅 Spark 指南。
Nessie 与 Iceberg
在大多数场景下,Nessie 对于 Iceberg 而言与其他 Catalog 无异:提供一组表的逻辑组织,并为事务提供原子性。然而,使用 Nessie 还能带来其他有趣的可能。当 Nessie 与 Iceberg 一起使用时,每一个 Iceberg 事务都会成为一次 Nessie 提交。这一历史记录可以被列示、合并,或在分支间进行 cherry-pick。
松耦合事务
通过创建一个分支并在该分支上执行一组操作,你可以近似实现多表事务。你可以在新创建的分支上执行一系列提交,然后将它们原子地合并回主分支。这样,一系列相关的更改看起来会同时对主分支生效。虽然下游消费者会看到多个事务同时出现,但这并不是数据库层面真正的多表事务。它本质上是(用 git 的说法)对多次提交的快进合并,分支上的每个操作都是各自独立的事务和提交。这与真正的多表事务不同——后者的所有更改都包含在同一次提交中。这种机制允许多个应用参与修改同一个分支,并使这组分布式的事务能够同时呈现给下游用户。
实验
对表的更改可以在分支中进行测试,然后再合并回主分支。这在执行架构演进(schema evolution)或分区演进(partition evolution)等大型变更时尤其有用。可以在一个分支中执行分区演进,然后在合并之前对这项变更进行测试(例如性能基准测试)。这为执行在线表修改和测试提供了极大的灵活性,同时不会中断下游使用场景。如果变更有误或性能不佳,可以直接丢弃该分支而不进行合并。
更多使用场景
有关 Nessie 特性的更多说明,请参阅 Nessie 文档。
危险
使用 Nessie 时,Iceberg 中的常规表维护会变得复杂。在执行任何表维护之前,请先查阅管理服务。
示例
有关 Nessie 与 Iceberg 结合使用的不同示例,请查看 Nessie Demos 仓库。
未来的改进
- Iceberg 多表事务:在同一事务中对多个 Iceberg 表进行更改、隔离级别等。
评论
登录后参与评论
KnowForge