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


##########
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:
   Looking up the current schema first fixes the renamed partition source at 
the current snapshot, but it moves the clash to `VERSION AS OF` reads. A time 
travel read's scan schema carries the snapshot's names, so the current name of 
an appended column can collide with a column the snapshot still has.
   
   I tried a table partitioned by `bucket(4, p)` that also has a column `c`. 
After taking a snapshot, I dropped `c` and renamed `p` to `c`. `SELECT id, c 
FROM t VERSION AS OF <snapshot>` then fails at planning with 
`ValidationException: Invalid schema: multiple fields for name c: 3 and 2`, 
with the config on or off, because `p` is appended under its current name next 
to the snapshot's `c`. On `main` and `branch-1.1` the same query matches Spark, 
because `p` comes back under its oldest name. Swapping two column names when 
one of them is a bucket partition source hits the same clash with pruning on 
(`multiple fields for name b`). `main` returns the other column's values there, 
which is the bug in my other comment.
   
   So neither order works for every case. What do you think about resolving the 
id from the schema the scan was planned with first? That is the snapshot's 
schema for `VERSION AS OF` and the current schema otherwise, and its names 
can't collide with the scan schema's. Skipping a name the base schema already 
has and trying the next schema would also work. Could the time travel test 
cover the drop-and-rename case too?



##########
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:
   The pruned scan schema also fixes a wrong-results bug in time travel, and 
this test could pin it down. On `main` and `branch-1.1`, a `VERSION AS OF` read 
of a column renamed after the snapshot returns NULL from the native scan. The 
task reads with the current table schema, so the column comes back under its 
new name and no longer matches the name Spark asks for. When two columns swap 
names, each returns the other's values. With pruning on, both match Spark, for 
tasks with deletes too.
   
   Could this test rename a column after `snapshotWithDeletes` and read it 
under the old name, plus a swap, so the fix stays covered? I filed #6822 for 
what's left, which is `branch-1.1`, the config turned off, and tasks that keep 
the full schema.



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