Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion docs/source/contributor-guide/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)).
Expand Down
8 changes: 7 additions & 1 deletion docs/source/user-guide/latest/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.*` —
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,7 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] {
requireFormatVersionAtMostTwo,
requireNoUuidColumns,
requireNoFloatingPointPartitionField,
requireNoVoidFieldWithDroppedSource,
requireNoEncryptionPrefix,
requireNoBloomFilterColumnsEnabled,
requireRowGroupCheckMinRecordCountAtDefault,
Expand Down Expand Up @@ -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))
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
}
Expand Down
Loading