From bc4fe7147abf6322397e0b8fa78dd4c2c50da8a3 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Mon, 5 Oct 2026 09:47:05 -0600 Subject: [PATCH] fix: decline native Iceberg writes through a void partition field whose source column was dropped A format-version-1 table keeps a dropped partition field as a void transform, and its source column can be dropped afterwards. The native writer failed every task on a spec that mixes such a field with a live one, because it resolves the spec's partition type against the schema. iceberg-java cannot write through that spec either, so the gate now declines it and the write fails with iceberg-java's own error. A spec whose fields are all void still writes natively, unpartitioned. --- .../contributor-guide/iceberg-writes.md | 3 +- .../user-guide/latest/iceberg-writes.md | 8 ++- .../comet/iceberg/IcebergReflection.scala | 31 +++++++++++ .../operator/CometIcebergNativeWrite.scala | 24 ++++++++ .../CometIcebergWriteDetectionSuite.scala | 55 ++++++++++++++++++- 5 files changed, 117 insertions(+), 4 deletions(-) diff --git a/docs/source/contributor-guide/iceberg-writes.md b/docs/source/contributor-guide/iceberg-writes.md index 98901b474e4..536813acb4f 100644 --- a/docs/source/contributor-guide/iceberg-writes.md +++ b/docs/source/contributor-guide/iceberg-writes.md @@ -401,7 +401,8 @@ Each of these has caused a bug on this path: ([#6145](https://github.com/apache/datafusion-comet/issues/6145)). - **Partition evolution leaves `void` fields behind.** A v1 spec keeps a dropped partition field as a `void` transform, whose source column may later be dropped from the schema. Resolving the spec - against the schema then fails + against the schema then fails. An all-`void` spec is written unpartitioned, and a spec that mixes + such a field with a live one is declined by the gate, since iceberg-java cannot write it either ([#5691](https://github.com/apache/datafusion-comet/issues/5691), [#5693](https://github.com/apache/datafusion-comet/issues/5693), [#6141](https://github.com/apache/datafusion-comet/issues/6141)). diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index 0b340dbbe27..a4e57ce350a 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -182,7 +182,7 @@ A write is eligible only when ALL of the following hold: | `write.target-file-size-bytes` | any value (the two writers can choose different roll points; see accepted divergences) | | data location URI scheme | `file`, `memory`, `s3`, `s3a`, `gs`, matched case-sensitively (`S3://` falls back). `s3`, `s3a` and `gs` need a bucket in the authority (`s3://bucket/...`), so a hostless form such as `s3:/bucket/key` falls back. `gs` only when the `FileIO` opening the data location is a `GCSFileIO`; see below | | resolved `table.locationProvider()` | Iceberg's built-in `DefaultLocationProvider` | -| partition spec | any, except an identity partition on a `float` or `double` column (see below) | +| partition spec | any, except an identity partition on a `float` or `double` column, or a `void` field whose source column was dropped beside a live field (see below) | | column types | any except `uuid` (Spark plans it as a string; no Arrow cast reaches `fixed(16)`) | Within the namespaces that shape data-file bytes — `write.parquet.*` and `parquet.*` — @@ -218,6 +218,12 @@ that prunes on the other value's partition would then miss rows. The fall-back s iceberg-rust distinguishes the two values ([apache/iceberg-rust#3325](https://github.com/apache/iceberg-rust/issues/3325)). +A format-version-1 table keeps a dropped partition field as a `void` field, and its source column +can be dropped afterwards. A spec that mixes such a field with a live one falls back: iceberg-java +cannot write through it either, so the write fails with iceberg-java's own error rather than in +the native writer ([#6141](https://github.com/apache/datafusion-comet/issues/6141)). A spec whose +fields are all `void` writes unpartitioned and stays eligible. + Other `write.*` properties are intentionally not gated because they cannot make the native writer produce different data files: distribution and ordering settings shape the Spark plan identically on both paths, WAP / branch / snapshot properties act on the JVM committer, diff --git a/spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala b/spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala index 37a13c76742..4952a857b60 100644 --- a/spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala +++ b/spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala @@ -804,6 +804,37 @@ object IcebergReflection extends Logging { .toSeq } + /** + * The names of the `void` partition fields of `spec` whose source column is no longer in the + * spec's schema, when `spec` also has a field that is not `void`. A format-version-1 table + * keeps a dropped partition field as a `void` transform, and its source column can be dropped + * afterwards. Empty when every field is `void`, since such a spec writes unpartitioned. Throws + * on reflection failure, so the caller can fail closed. + */ + def voidFieldsWithDroppedSource(spec: Any): Seq[String] = { + import scala.jdk.CollectionConverters._ + val schema = getMethod(spec.getClass, "schema").invoke(spec) + val findField = getMethod(schema.getClass, "findField", classOf[Int]) + val fields = + getMethod(spec.getClass, "fields") + .invoke(spec) + .asInstanceOf[java.util.List[_]] + .asScala + .toSeq + def isVoid(field: Any): Boolean = + getMethod(field.getClass, "transform").invoke(field).toString == "void" + if (fields.forall(isVoid)) { + Seq.empty + } else { + fields + .filter { field => + val sourceId = getMethod(field.getClass, "sourceId").invoke(field).asInstanceOf[Int] + isVoid(field) && findField.invoke(schema, sourceId.asInstanceOf[Object]) == null + } + .map(field => getMethod(field.getClass, "name").invoke(field).asInstanceOf[String]) + } + } + /** * Gets the partition spec from an Iceberg table. */ diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala index 5d5510fb7a1..b9d9c684503 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala @@ -180,6 +180,7 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] { requireFormatVersionAtMostTwo, requireNoUuidColumns, requireNoFloatingPointPartitionField, + requireNoVoidFieldWithDroppedSource, requireNoEncryptionPrefix, requireNoBloomFilterColumnsEnabled, requireRowGroupCheckMinRecordCountAtDefault, @@ -295,6 +296,29 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] { } } + // A format-version-1 spec keeps a dropped partition field as a `void` transform, and its source + // column can be dropped afterwards. iceberg-java cannot write through a spec that mixes such a + // field with a live one, and the native writer fails resolving the spec's partition type, so + // decline and let the write fail the way iceberg-java fails it. An all-`void` spec stays + // eligible, since the native writer writes it unpartitioned. + // https://github.com/apache/datafusion-comet/issues/6141 + private val requireNoVoidFieldWithDroppedSource: TriggerRule = ctx => + IcebergReflection + .getOutputSpecIdFromSparkWrite(ctx.sparkWrite) + .flatMap(IcebergReflection.getPartitionSpecById(ctx.table, _)) match { + case None => Some("could not resolve the output partition spec for void field checking") + case Some(spec) => + try { + IcebergReflection.voidFieldsWithDroppedSource(spec).headOption.map { name => + s"partition field $name is a void transform whose source column was dropped, " + + "beside a live partition field" + } + } catch { + case e: Exception => + Some(s"could not inspect the output partition spec: ${e.getMessage}") + } + } + private val requireNoEncryptionPrefix: TriggerRule = ctx => ctx.properties.keys .find(_.startsWith(EncryptionPropertyPrefix)) diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala index bed4e85d552..5b819a0ce86 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala @@ -1030,6 +1030,55 @@ class CometIcebergWriteDetectionSuite extends CometTestBase with CometIcebergTes } } + // A format-version-1 spec keeps a dropped partition field as a `void` transform, and the + // field's source column can be dropped afterwards. Neither writer can write through a spec that + // mixes such a field with a live one: iceberg-java fails building the partition key, and the + // native writer fails resolving the spec's partition type. Declining keeps the failure + // iceberg-java's own. The insert fails either way, so only the gate's decision is checked, on + // an insert that is planned but not run. + // https://github.com/apache/datafusion-comet/issues/6141 + test("fall-back: void partition field whose source column was dropped, beside a live field") { + withDetectionCatalog { dir => + createTable( + dir, + "void_dropped", + partitionSpec = "PARTITIONED BY (region, id)", + properties = Some("'format-version'='1'")) + // Loaded afresh for each change, as each change commits a new metadata version. + def table: org.apache.iceberg.Table = + loadIcebergTable(spark, catalog, ns, "void_dropped") + .asInstanceOf[org.apache.iceberg.Table] + table.updateSpec().removeField("id").commit() + table.updateSchema().deleteColumn("id").commit() + spark.sql(s"REFRESH TABLE $catalog.$ns.void_dropped") + val writeExec = planInsertWriteExec(s"$catalog.$ns.void_dropped", values = "('us', 1.0)") + assertUnsupportedContains(writeExec, "void_dropped", "void", "source column", "dropped") + } + } + + test("Compatible when every partition field is void, even with its source column dropped") { + // Such a spec writes unpartitioned, so the native writer needs no source column for it. On + // Iceberg 1.8 neither writer can write to this table, since iceberg-java cannot bind the + // original spec on the executors once its source column is gone, so only the gate's decision + // is checked, on an insert that is planned but not run. + withDetectionCatalog { dir => + createTable( + dir, + "all_void_dropped", + partitionSpec = "PARTITIONED BY (region)", + properties = Some("'format-version'='1'")) + def table: org.apache.iceberg.Table = + loadIcebergTable(spark, catalog, ns, "all_void_dropped") + .asInstanceOf[org.apache.iceberg.Table] + table.updateSpec().removeField("region").commit() + table.updateSchema().deleteColumn("region").commit() + spark.sql(s"REFRESH TABLE $catalog.$ns.all_void_dropped") + val writeExec = planInsertWriteExec(s"$catalog.$ns.all_void_dropped", values = "(1, 1.0)") + val support = CometIcebergNativeWrite.getSupportLevel(writeExec) + assert(support.isInstanceOf[Compatible], s"expected Compatible, got $support") + } + } + test("fall-back: uuid column in the write schema") { withDetectionCatalog { dir => // Spark DDL cannot declare `uuid`, so evolve the schema through the Iceberg API. Spark @@ -1088,9 +1137,11 @@ class CometIcebergWriteDetectionSuite extends CometTestBase with CometIcebergTes * `CommandExecutionMode.SKIP` keeps `QueryExecution` from eagerly running the write command, so * a data location no filesystem on this classpath can reach never triggers a write. */ - private def planInsertWriteExec(qualifiedTable: String): IcebergWriteExec = { + private def planInsertWriteExec( + qualifiedTable: String, + values: String = "(1, 'us', 1.0)"): IcebergWriteExec = { val plan = - spark.sessionState.sqlParser.parsePlan(s"INSERT INTO $qualifiedTable VALUES (1, 'us', 1.0)") + spark.sessionState.sqlParser.parsePlan(s"INSERT INTO $qualifiedTable VALUES $values") findWriteExecOrFail( spark.sessionState.executePlan(plan, CommandExecutionMode.SKIP).executedPlan) }