comphead commented on code in PR #6725:
URL: https://github.com/apache/datafusion-comet/pull/6725#discussion_r4238758740


##########
spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala:
##########
@@ -741,12 +741,14 @@ 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)
+      // table.schemas() lists the oldest schema first, so a column the table 
still has is looked
+      // up in the current schema before it, to keep its current name and 
type. An older name
+      // could clash with a column that took it over.
+      val schemas = getMethod(table.getClass, "schema").invoke(table) +: 
getAllSchemas(table)

Review Comment:
   You're right, thanks for the repro. Iceberg's `Schema` constructor indexes 
the names and rejects the second `c`, and the config-off path goes through the 
same `schemaWithRequiredFields`, so the switch doesn't help. In d88698c442 the 
missing fields all come from one schema: the current one, or the newest in the 
history that has them all under names the task schema does not use. That covers 
your drop-and-rename and swap cases, and also the two-source case sunchao 
found, where picking a name per field still clashed. It keeps current names and 
types when they don't clash, and needs no extra reflection to find the scan's 
snapshot. The time travel test covers the drop-and-rename case on the 
partitioned table: after `c` is dropped and `p` is renamed to `c`, `SELECT id, 
c, s.a ... VERSION AS OF` reads the snapshot's `c`, and `p` is appended under 
its older name.
   



##########
spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala:
##########
@@ -2284,6 +2285,256 @@ class CometIcebergNativeSuite
     }
   }
 
+  // Spark's nested schema pruning reaches the native scan through the scan 
schema, so only the
+  // nested fields a query uses are read and the wide `pad` strings beside 
them are skipped.
+  // sql-tests/iceberg/nested_schema_pruning.sql checks the results. One data 
file holds every
+  // row, so each `pad` column chunk is larger than iceberg-rust's 1 MiB read 
coalescing, which
+  // would otherwise merge the reads of the kept chunks across the skipped 
ones.
+  test("nested schema pruning reads only the nested fields the query uses") {
+    assume(icebergAvailable, "Iceberg not available in classpath")
+
+    withTempIcebergDir { warehouseDir =>
+      withSQLConf(
+        "spark.sql.catalog.test_cat" -> 
"org.apache.iceberg.spark.SparkCatalog",
+        "spark.sql.catalog.test_cat.type" -> "hadoop",
+        "spark.sql.catalog.test_cat.warehouse" -> warehouseDir.getAbsolutePath,
+        CometConf.COMET_ENABLED.key -> "true",
+        CometConf.COMET_EXEC_ENABLED.key -> "true",
+        CometConf.COMET_ICEBERG_NATIVE_ENABLED.key -> "true") {
+
+        val table = "test_cat.db.nested_pruning"
+        spark.sql(s"""
+          CREATE TABLE $table (
+            id INT,
+            s STRUCT<a: INT, pad: STRING>,
+            items ARRAY<STRUCT<x: INT, pad: STRING>>,
+            m MAP<STRING, STRUCT<v: INT, pad: STRING>>
+          ) USING iceberg
+        """)
+        spark.sql(s"""
+          INSERT INTO $table
+          SELECT
+            CAST(id AS INT),
+            named_struct('a', CAST(id AS INT), 'pad', pad),
+            array(named_struct('x', CAST(id AS INT), 'pad', pad)),
+            map('k', named_struct('v', CAST(id AS INT), 'pad', pad))
+          FROM (
+            SELECT id, concat_ws('', transform(array('a', 'b', 'c', 'd'),
+              salt -> sha2(concat(CAST(id AS STRING), salt), 256))) AS pad
+            FROM range(0, 20000, 1, 1))
+        """)
+
+        def bytesScanned(pruneNestedFields: Boolean): Long = {
+          var bytes = 0L
+          withSQLConf(
+            CometConf.COMET_ICEBERG_NESTED_SCHEMA_PRUNING_ENABLED.key ->
+              pruneNestedFields.toString) {
+            val df = spark.sql(s"SELECT sum(s.a), count(items.x), 
count(m['k'].v) FROM $table")
+            df.collect()
+            val scans = 
collectIcebergNativeScans(df.queryExecution.executedPlan)
+            assert(scans.length == 1, s"expected one native scan, got 
${scans.length}")
+            bytes = scans.head.metrics("bytes_scanned").value
+          }
+          bytes
+        }
+        val prunedBytes = bytesScanned(pruneNestedFields = true)
+        val fullBytes = bytesScanned(pruneNestedFields = false)
+        assert(
+          prunedBytes * 4 < fullBytes,
+          s"pruned read should skip the pad fields: pruned=$prunedBytes, 
full=$fullBytes")
+
+        spark.sql(s"DROP TABLE $table")
+      }
+    }
+  }
+
+  // A pruned task schema still needs the columns iceberg-rust uses beyond the 
projection: the
+  // partition source and the equality-delete key when the query projects 
neither. Tasks with
+  // deletes read with the pruned schema too, and the tasks of every partition 
share the one schema
+  // that appending the partition source builds.
+  test("nested schema pruning with deletes, partitions, and time travel") {
+    assume(icebergAvailable, "Iceberg not available in classpath")
+
+    withTempIcebergDir { warehouseDir =>
+      withSQLConf(
+        "spark.sql.catalog.test_cat" -> 
"org.apache.iceberg.spark.SparkCatalog",
+        "spark.sql.catalog.test_cat.type" -> "hadoop",
+        "spark.sql.catalog.test_cat.warehouse" -> warehouseDir.getAbsolutePath,
+        CometConf.COMET_ENABLED.key -> "true",
+        CometConf.COMET_EXEC_ENABLED.key -> "true",
+        CometConf.COMET_ICEBERG_NATIVE_ENABLED.key -> "true") {
+
+        // Also checks that no task, tasks with deletes included, reads the 
pruned `pad` fields.
+        // With `oneSchema`, every task needs the same columns, so they share 
one task schema.
+        def checkPrunedNativeScan(query: String, oneSchema: Boolean = false): 
Unit = {
+          val (_, cometPlan) = checkSparkAnswer(query)
+          val scans = collectIcebergNativeScans(cometPlan)
+          assert(scans.length == 1, s"expected one native scan, got 
${scans.length}\n$cometPlan")
+          // Planning commonData leaks manifest streams on Iceberg versions 
before 1.8.0.
+          if (icebergVersionAtLeast(1, 8)) {
+            val schemas = OperatorOuterClass.IcebergScanCommon
+              .parseFrom(scans.head.commonData)
+              .getSchemaPoolList
+              .asScala
+            assert(
+              schemas.forall(!_.contains("\"pad\"")),
+              s"$query: pruned field in a task 
schema:\n${schemas.mkString("\n")}")
+            assert(
+              !oneSchema || schemas.length == 1,
+              s"$query: expected one task schema, 
got:\n${schemas.mkString("\n")}")
+          }
+        }
+        def latestSnapshotId(table: String): Long = spark
+          .sql(s"SELECT snapshot_id FROM $table.snapshots ORDER BY 
committed_at DESC LIMIT 1")
+          .collect()(0)
+          .getLong(0)
+
+        val morProperties = """
+          TBLPROPERTIES (
+            'format-version' = '2',
+            'write.delete.mode' = 'merge-on-read',
+            'write.update.mode' = 'merge-on-read',
+            'write.merge.mode' = 'merge-on-read')
+        """
+        val rows = """
+          SELECT CAST(id AS INT) AS id, IF(id % 2 = 0, 'even', 'odd') AS p,
+            named_struct('a', CAST(id AS INT), 'pad', repeat('x', 100)) AS s
+          FROM range(200)
+        """
+
+        val mor = "test_cat.db.nested_pruning_mor"
+        spark.sql(s"""
+          CREATE TABLE $mor (id INT, c STRING, s STRUCT<a: INT, pad: STRING>)
+          USING iceberg $morProperties
+        """)
+        spark.sql(s"INSERT INTO $mor SELECT id, p, s FROM ($rows)")
+        val snapshotBeforeDeletes = latestSnapshotId(mor)
+        spark.sql(s"DELETE FROM $mor WHERE id % 10 = 0")
+        val snapshotWithDeletes = latestSnapshotId(mor)
+        commitEqualityDelete("test_cat", "db", "nested_pruning_mor", "id", 7, 
warehouseDir)
+        spark.sql(s"ALTER TABLE $mor DROP COLUMN c")
+        checkPrunedNativeScan(s"SELECT id, s.a FROM $mor ORDER BY id")
+        // The equality-delete key `id` is not projected. Iceberg gives the 
delete only to the data
+        // file whose `id` range holds 7, so only that file's task appends 
`id`.
+        checkPrunedNativeScan(s"SELECT s.a FROM $mor ORDER BY s.a")
+        checkPrunedNativeScan(
+          s"SELECT id, s.a FROM $mor VERSION AS OF $snapshotBeforeDeletes 
ORDER BY id")
+        // `c` was dropped after this snapshot, so the current table schema 
lacks it, and these
+        // tasks with deletes must read with the scan schema.
+        checkPrunedNativeScan(
+          s"SELECT id, c, s.a FROM $mor VERSION AS OF $snapshotWithDeletes 
ORDER BY id")

Review Comment:
   Thanks for finding this and filing #6822. The test covers it since 
335e7e0caf: after `snapshotWithDeletes`, `r` is renamed to `r2` and `x` and `y` 
swap names, next to the dropped `c`, and `SELECT id, c, r, x, y, s.a ... 
VERSION AS OF` reads them under the snapshot's names, with position deletes.
   



##########
spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala:
##########
@@ -741,28 +743,68 @@ 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)
-      resolvedFields.foreach(newColumns.add)
+      // A field is appended under a name it has had that the task schema does 
not use yet. The
+      // current schema comes first, to keep a live column's current name and 
type, and then
+      // table.schemas(), oldest first. A VERSION AS OF scan schema carries 
the snapshot's names,
+      // so a current name can clash with a column the snapshot still has.
+      val schemas = getMethod(table.getClass, "schema").invoke(table) +: 
getAllSchemas(table)
+      val names = 
scala.collection.mutable.Set(existing.asScala.map(fieldName).toSeq: _*)
+      missingIds.foreach { id =>
+        val (field, name) = schemas.iterator
+          .flatMap(findFieldObject(_, id))
+          .map(f => (f, fieldName(f)))
+          .find { case (_, n) => !names.contains(n) }

Review Comment:
   Thanks, this holds from the code. The partition source ids come back in spec 
order, `[2, 3]`, id 2 takes its current name `q`, and id 3's names `c` and `q` 
are then both taken. Fixed in d88698c442: `schemaWithRequiredFields` now takes 
all the missing ids from one schema, the current one or the newest in the 
history that has them all under names the task schema does not use. One schema 
cannot give two fields the same name, so the order of the ids no longer 
matters, and newest first keeps a promoted column's widest type. I searched the 
history rather than looking up the scan's snapshot schema, because it needs no 
further reflection into Iceberg's scan internals, and whenever the snapshot's 
schema has all the ids it qualifies too. The time travel test now has your 
case, with both sources bucketed and `c` read through `VERSION AS OF`.
   



-- 
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