Databend 测试基建中的 Iceberg Driver:用 Iceberg Java 高层 API 直接生成 Equality Delete 测试数据
【免费下载链接】databendData Agent Ready Warehouse : One for Analytics, Search, AI, Python Sandbox. — rebuilt from scratch. Unified architecture on your S3.项目地址: https://gitcode.com/GitHub_Trending/da/databend
本指南讲解 Databend 仓库tests/sqllogictests/scripts/iceberg-driver目录下的 Iceberg Driver 工具:它绕过 SQL 层,直接调用 Iceberg Java 高层 API 向 Iceberg 表写入 Equality Delete 数据文件,用于验证 Databend 对 Iceberg 外部目录中 Merge-on-Read 删除语义的兼容性。读完本文,你将掌握该工具的整体架构、Maven 工程依赖、Driver 主程序的完整执行流程,以及它在 Docker 测试环境(Iceberg REST Catalog + RustFS S3)中的运行方式,并了解 Databend 侧如何通过CREATE CATALOG ... TYPE = ICEBERG挂载同一份数据。
一、为什么需要一个独立的 Iceberg Driver
Iceberg 表的删除分为Position Delete(位置删除,按数据文件内的行位置标记删除)与Equality Delete(等值删除,按列值匹配删除)。其中 Equality Delete 在业界 SQL 引擎中很少直接暴露,很多引擎只生成 Position Delete,因此用常规 SQL 手段很难构造出"同时包含两类删除文件、且相互叠加"的复杂表状态。
Databend 的 Iceberg 外部目录测试需要这类数据来验证读取语义,因此仓库在 tests/sqllogictests/scripts/iceberg-driver/README.md 中明确说明了该目录的定位:
Use some high-level APIs directly in Iceberg Java to generate data — For: Equality Delete
即:在 Iceberg Java 中直接使用高层 API 生成数据,目标是构造 Equality Delete 场景。这与 Databend 的 Iceberg 测试套件 tests/suites/3_stateful_iceberg 相互配合,为状态化(stateful)测试提供真实、可控的数据源。
二、整体架构:Docker 编排下的四服务测试环境
Iceberg Driver 不是一个独立运行的普通程序,而是被编排进一套完整的测试环境中。仓库中的 docker-compose-iceberg-tpch.yml 定义了 4 个服务:
| 服务 | 镜像 | 作用 |
|---|---|---|
rest | apache/iceberg-rest-fixture:1.10.0 | 提供 Iceberg REST Catalog,监听127.0.0.1:8181 |
rustfs | rustfs/rustfs:1.0.0-alpha.91 | S3 兼容对象存储,监听127.0.0.1:9002,作为 Iceberg 的 warehouse |
mc | minio/mc | 初始化 bucketiceberg-tpch并清理旧数据 |
iceberg-driver | 本地构建(build: context: ./iceberg-driver) | 本文主角,负责生成测试数据 |
REST Catalog 的关键环境变量如下(均与 Driver.java 中的连接参数一一对应):
environment: - AWS_ACCESS_KEY_ID=admin - AWS_SECRET_ACCESS_KEY=password - AWS_REGION=us-east-1 - CATALOG_WAREHOUSE=s3://iceberg-tpch/ - CATALOG_IO__IMPL=org.apache.iceberg.aws.s3.S3FileIO - CATALOG_S3_ENDPOINT=http://127.0.0.1:9002 - CATALOG_S3_ACCESS__KEY__ID=admin - CATALOG_S3_SECRET__ACCESS__KEY=password - CATALOG_CLIENT_REGION=us-east-1注意所有服务都使用了network_mode: "host"(主机网络模式),因此各组件之间通过127.0.0.1直接互通。mc服务负责等待 RustFS 就绪后执行mc mb rustfs/iceberg-tpch创建 bucket,并设置 public 策略;iceberg-driver通过depends_on: mc确保 bucket 已就绪后才启动。
三、Maven 工程:依赖与打包方式
工程定义在 pom.xml 中,groupId为org.example,artifactId为iceberg-equality-delete,主类为org.example.Driver。
关键依赖版本(<properties>段):
- Java 编译目标:11(
maven.compiler.source/target) - Scala 二进制版本:
2.13,Spark 二进制版本:3.3 - Spark:
3.3.3(spark-core_2.13与spark-sql_2.13) - Hadoop:
3.3.0 - Iceberg:
1.5.1 - AWS SDK:
2.20.80,Jackson:2.13.4.2
引入的核心依赖:
<dependency> <groupId>org.apache.iceberg</groupId> <artifactId>iceberg-aws-bundle</artifactId> <version>1.5.1</version> </dependency> <dependency> <groupId>org.apache.iceberg</groupId> <artifactId>iceberg-spark-runtime-${spark.binary.version}_${scala.binary.version}</artifactId> <version>${iceberg.veriosn}</version> </dependency> <dependency> <groupId>org.apache.iceberg</groupId> <artifactId>iceberg-spark-extensions-${spark.binary.version}_${scala.binary.version}</artifactId> <version>${iceberg.veriosn}</version> </dependency>打包使用maven-assembly-plugin生成fat jar(jar-with-dependencies,且appendAssemblyId=false使最终产物直接命名为app.jar形式的可执行 jar),并在 manifest 中声明mainClass=org.example.Driver。这意味着java -jar app.jar即可直接运行,无需额外提供 classpath。
Dockerfile(Dockerfile)采用两阶段构建:先在maven:3.9.5-eclipse-temurin-11中执行mvn -B -ntp clean package -DskipTests(网络不稳定时最多重试 5 次,每次递增等待 5 秒),再把产物拷贝到eclipse-temurin:11-jdk运行时镜像。启动命令为:
CMD ["java", "--add-exports=java.base/sun.nio.ch=ALL-UNNAMED", "-jar", "app.jar"]--add-exports用于在 JDK 11 模块化环境下开放sun.nio.ch内部包,满足 Spark 运行时的反射访问需求。
四、Driver 主程序逐步拆解
核心实现位于 Driver.java,它分四个阶段工作:Spark 建表 → SQL 删除 → Java API 写 Equality Delete → 回插数据。
4.1 通过 Spark 配置 Iceberg REST Catalog
程序首先构建一个本地 SparkSession,将iceberg注册为 REST 类型的 Spark Catalog:
SparkSession spark = SparkSession.builder() .appName("CSV to Iceberg REST Catalog").master("local[*]") .config("spark.sql.catalog.iceberg", "org.apache.iceberg.spark.SparkCatalog") .config("spark.sql.catalog.iceberg.type", "rest") .config("spark.sql.catalog.iceberg.uri", "http://127.0.0.1:8181") .config("spark.sql.catalog.iceberg.io-impl", "org.apache.iceberg.aws.s3.S3FileIO") .config("spark.sql.catalog.iceberg.warehouse", "s3://iceberg-tpch/") .config("spark.sql.catalog.iceberg.s3.access-key-id", "admin") .config("spark.sql.catalog.iceberg.s3.secret-access-key", "password") .config("spark.sql.catalog.iceberg.s3.path-style-access", "true") .config("spark.sql.catalog.iceberg.s3.endpoint", "http://127.0.0.1:9002") .config("spark.sql.catalog.iceberg.client.region", "us-east-1") .config("spark.jars.packages", "org.apache.iceberg:iceberg-aws-bundle:1.6.1,org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.6.1") .getOrCreate();配置要点说明:
spark.jars.packages在运行时按需拉取 Iceberg 1.6.1 的 AWS bundle 与 Spark 3.5 运行时(工程编译期使用 1.5.1,运行期使用 1.6.1,二者均为 Apache Iceberg 官方构件)。s3.path-style-access=true表示使用 path-style 寻址访问本地 S3 兼容存储(RustFS),s3.endpoint=http://127.0.0.1:9002指向 RustFS 服务端口。
4.2 创建 Merge-on-Read 表并准备基础数据
Driver 在iceberg.test命名空间下创建目标表,并通过表属性强制Merge-on-Read删除模式与Format v2:
spark.sql("CREATE OR REPLACE TABLE iceberg.test.test_merge_on_read_deletes (\n" + " dt date,\n" + " number integer,\n" + " letter string\n" + ")\n" + "USING iceberg\n" + "TBLPROPERTIES (\n" + " 'write.delete.mode'='merge-on-read',\n" + " 'write.update.mode'='merge-on-read',\n" + " 'write.merge.mode'='merge-on-read',\n" + " 'format-version'='2'\n" + ");");随后插入 12 行数据(2023-03-01到2023-03-12,number为 1 到 12,letter为 'a' 到 'l')。
4.3 SQL 层删除:产生 Position Delete
接着执行一条常规 SQL DELETE:
spark.sql("DELETE FROM iceberg.test.test_merge_on_read_deletes WHERE number > 4 AND number < 7");该语句删除number = 5与number = 6两行。由于表开启了write.delete.mode = merge-on-read,这条 SQL 会生成Position Delete 文件——这是"引擎默认路径"产生的删除数据。
4.4 Java 高层 API:直接写 Equality Delete 文件
这是本工具的核心亮点。Driver 绕过 Spark SQL,直接使用 Iceberg 的高层 API(RESTCatalog+RowDelta+FileMetadata.deleteFileBuilder)向表中注入 Equality Delete。
首先用与 Spark 相同的连接参数初始化RESTCatalog并加载表:
Map<String, String> properties = new HashMap<>(); properties.put(CatalogProperties.URI, "http://127.0.0.1:8181"); properties.put(CatalogProperties.WAREHOUSE_LOCATION, "s3://iceberg-tpch/"); properties.put("io-impl", "org.apache.iceberg.aws.s3.S3FileIO"); properties.put("s3.access-key-id", "admin"); properties.put("s3.secret-access-key", "password"); properties.put("s3.endpoint", "http://127.0.0.1:9002"); properties.put("s3.path-style-access", "true"); properties.put("s3.region", "us-east-1"); RESTCatalog restCatalog = new RESTCatalog(); restCatalog.initialize("rest", properties); TableIdentifier tableId = TableIdentifier.of("test", "test_merge_on_read_deletes"); Table table = restCatalog.loadTable(tableId);然后构造两个删除记录,分别针对单列等值语义:
// 删除 number = 3 的行 Schema number = table.schema().select("number"); Record deleteRecord0 = GenericRecord.create(number); deleteRecord0.setField("number", 3); // 删除 letter = 'k' 的行(即 number = 11 那一行) Schema letter = table.schema().select("letter"); Record deleteRecord1 = GenericRecord.create(letter); deleteRecord1.setField("letter", "k");DeleteRecord是文件内定义的一个辅助内部类,同时持有schema、record与fieldName,并提供fieldId()方法:
int fieldId() { return this.schema.findField(this.fieldName).fieldId(); }fieldId()通过Schema.findField拿到目标字段,再取其 Iceberg 内部的fieldId()(Format v2 下 Equality Delete 按字段 ID 而非字段名定位),这是构造合法 Equality Delete 文件的必要条件。
对每一个删除记录,Driver 执行如下四步:
第一步:用 Parquet 写删除数据文件
OutputFile outputFile = table.io().newOutputFile(deleteFilePath); FileAppender<Record> appender = Parquet.write(outputFile) .schema(deleteRecords.get(i).schema) .createWriterFunc(GenericParquetWriter::buildWriter) .build(); appender.add(deleteRecords.get(i).record); appender.close();删除文件路径形如s3://iceberg-tpch/test/test_merge_on_read_deletes/data/equality-delete-file-{i}.parquet,通过GenericParquetWriter::buildWriter以 Generic Record 方式落盘。删除文件的 schema 是投影后的单字段 schema(table.schema().select("number")/select("letter")),这正是 Equality Delete 文件的数据形态:只包含用于匹配的等值列。
第二步:构造 DeleteFile 元数据
DeleteFile deleteFile = FileMetadata.deleteFileBuilder(spec) .ofEqualityDeletes(deleteRecords.get(i).fieldId()) .withFormat(FileFormat.PARQUET) .withPath(deleteFilePath) .withPartition(partitionData) .withFileSizeInBytes(appender.length()) .withMetrics(appender.metrics()) .withSplitOffsets(appender.splitOffsets()) .withSortOrder(table.sortOrder()) .build();ofEqualityDeletes(fieldId)明确声明该文件为等值删除;withMetrics/withSplitOffsets直接复用FileAppender写入过程中统计到的列指标与行组偏移信息,保证元数据与文件内容一致。
第三步:加入 RowDelta 事务
rowDelta.addDeletes(deleteFile);第四步:提交
rowDelta.commit();RowDelta是 Iceberg 中"添加删除文件 + 重写数据文件"的原子变更接口,此处仅添加删除文件,未添加数据文件,因此最终表快照中会同时存在 Position Delete 与两类单列 Equality Delete。
4.5 回插被删数据,构造"删除后又插入"的复杂状态
提交删除后,Driver 再补一条 INSERT,把number = 6重新插回去:
spark.sql("INSERT INTO iceberg.test.test_merge_on_read_deletes VALUES (CAST('2023-03-30' AS date), 6, 'z');");这样最终表的删除语义变得非常"刁钻":既有按number的等值删除(3),又有按letter的等值删除('k'),还有先按 range 删除number=6、随后又插入新行number=6, letter='z'的新旧版本纠缠。Databend 在读取该表时必须同时正确处理 Position Delete、两列 Equality Delete,以及"新插入行不受旧删除文件影响"的语义,这正是该工具存在的价值。
五、与 Databend 侧 Iceberg Catalog 的衔接
数据生成后,Databend 通过 Iceberg Catalog 直接读取这份数据。仓库中 tests/suites/3_stateful_iceberg/00_rest/00_0000_create_and_show.sh 展示了 Databend 侧的挂载方式:
cat <<EOF | bendsql_connect_root CREATE CATALOG iceberg_rest TYPE = ICEBERG CONNECTION = ( TYPE = 'rest' ADDRESS = 'http://localhost:8181' warehouse = 's3://warehouse/demo/' "s3.endpoint" = 'http://localhost:9000' "s3.access-key-id" = 'admin' "s3.secret-access-key" = 'password' "s3.region" = 'us-east-1' ); EOF创建完成后即可在 Databend 中USE CATALOG iceberg_rest;、SHOW TABLES、执行查询,测试脚本随后会依次验证建表、写入、读取、删表等完整流程。可以看到,Databend 侧声明 Catalog 的TYPE='rest'、ADDRESS、warehouse、s3.endpoint等参数,与 Driver 侧RESTCatalog的初始化属性保持同一套语义,保证了双方访问的是同一份元数据与同一批数据文件。
六、同类工具对照与使用小结
在scripts目录中还有两个基于 PySpark 的数据准备脚本与本 Driver 形成互补:
- prepare_iceberg_tpch_data.py:使用 Spark 4.0 + Iceberg 1.10.0,将 TPC-H 生成的 CSV 数据批量灌入
iceberg.tpch命名空间,为 TPC-H 查询测试提供基础数据。 - prepare_iceberg_test_data.py:构造常规的 Iceberg 测试表。
它们都复用同一套 REST Catalog + RustFS 连接参数(http://127.0.0.1:8181/s3://iceberg-tpch//127.0.0.1:9002),与iceberg-driver共享同一个 Docker 测试环境。区别在于:PySpark 脚本走的是标准 Spark SQL 路径,只能生成引擎支持的数据形态;而 Iceberg Driver 直接调用 Java 高层 API,能够精确构造 PySpark 脚本无法产生的Equality Delete 文件,从源码结构看,这正是它被独立成目录、单独编排进docker-compose-iceberg-tpch.yml的原因。
七、小结
iceberg-driver是 Databend Iceberg 兼容性测试链路中一个小巧但不可替代的环节:
- 它用Spark SQL 建表 + Java 高层 API 注入删除的组合,一次性构造出 Position Delete、单列 Equality Delete 与"删除后重插"三类状态叠加的 Merge-on-Read 表;
- 它的连接参数与 REST Catalog、RustFS、Databend 侧
CREATE CATALOG完全对齐,可无缝融入既有测试环境; - 其完整实现(Driver.java)、构建配置(pom.xml)、运行镜像(Dockerfile)与编排文件(docker-compose-iceberg-tpch.yml)都保留在仓库中,可作为"如何用 Iceberg Java 高层 API 生成特殊删除场景数据"的现成范例复用。
【免费下载链接】databendData Agent Ready Warehouse : One for Analytics, Search, AI, Python Sandbox. — rebuilt from scratch. Unified architecture on your S3.项目地址: https://gitcode.com/GitHub_Trending/da/databend
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考