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]

Reply via email to