sunchao commented on code in PR #6725:
URL: https://github.com/apache/datafusion-comet/pull/6725#discussion_r4238792307
##########
spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala:
##########
@@ -741,20 +743,30 @@ object IcebergReflection extends Logging {
logDebug(
s"Native Iceberg scan schema is missing field id(s)
${missingIds.mkString(",")}; " +
"resolving them from table schema history")
- val history = getAllSchemas(table)
- val resolvedFields = missingIds.map { id =>
- history.iterator
- .flatMap(s => findFieldObject(s, id))
- .toSeq
- .headOption
- .getOrElse(throw new IllegalStateException(
- s"Cannot resolve field id $id in table schema history"))
- }
val existing =
getMethod(baseSchema.getClass, "columns")
.invoke(baseSchema)
.asInstanceOf[java.util.List[_]]
- val newColumns = new java.util.ArrayList[Any](existing)
+ val existingNames = existing.asScala.map(fieldName).toSet
+ // Taking every field from one schema keeps their names apart, which a
name picked per field
+ // cannot promise once columns have been renamed into each other's
names. A VERSION AS OF
+ // scan schema carries the snapshot's names, so a schema also must not
name a field like a
+ // column `baseSchema` has. The current schema comes first, to keep
current names, and then
+ // the others newest first, so that a promoted column keeps its widest
type.
+ val schemas = getMethod(table.getClass, "schema").invoke(table) +:
+ getAllSchemas(table).reverse
+ val resolvedFields = schemas.iterator
+ .map(schema => missingIds.flatMap(findFieldObject(schema, _)))
+ .find { fields =>
+ val names = fields.map(fieldName)
+ fields.length == missingIds.length && names.distinct.length ==
names.length &&
Review Comment:
[P2] [P2] Allow required delete keys from different schema generations. On
Iceberg 1.11, create an unpartitioned table `t(id INT, p INT)` with rows
`(1,10),(2,20),(3,30)`, commit an equality delete for `p=10`, drop `p`, add `q
INT`, and commit an equality delete for `q=99`. `SELECT id FROM t ORDER BY id`
should return `2,3`, as Spark does. Both delete files apply to the original
data file, so the default pruned path must append IDs `[2,3]` to the `id`
projection. No historical schema contains both keys, and this condition rejects
every candidate, throwing `Cannot resolve field ids 2,3 from one schema...`.
The base successfully augments the task schema, and disabling pruning also
succeeds at this helper boundary for this case. This breaks valid reads even
without nested columns. Preserve collision-safe resolution without requiring
all missing IDs to coexist in one schema, or fall back during planning.
Evidence: Reproduced with genuine Parquet data/delete files on Spark 4.1.3
and Iceberg 1.11.0. Spark returned `ArraySeq(2, 3)` and `planFiles()` attached
equality IDs `List(2, 3)` to the data task. Compiled the exact d88698c442
helper: the pruned call throws the stated exception, while the base helper
succeeds with fields `id,q,p`. Verified that the base helper equals
branch-1.1's implementation. Sources and logs:
`/tmp/pr6725-d88698-validation-ca4urq8k/SparkProbe.scala`, `Probe.scala`,
`spark-oracle.log`, and `evidence.json`. This validates serialization failure
directly, without claiming an end-to-end Comet run.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]