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


##########
spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala:
##########
@@ -2284,6 +2285,301 @@ 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. NULL
+  // structs, lists, maps, and list elements check that validity is rebuilt 
from the pruned leaves.
+  // 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, inner: STRUCT<b: 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),
+            IF(id % 7 = 0, NULL, named_struct(
+              'a', IF(id % 5 = 0, NULL, CAST(id AS INT)),
+              'pad', pad,
+              'inner', named_struct('b', CAST(id * 2 AS INT), 'pad', pad))),
+            IF(id % 7 = 0, NULL, array(
+              named_struct('x', CAST(id AS INT), 'pad', pad),
+              IF(id % 5 = 0, NULL, named_struct('x', CAST(-id AS INT), 'pad', 
pad)))),
+            IF(id % 7 = 0, NULL, 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))
+        """)
+
+        Seq(
+          s"SELECT id, s.a FROM $table ORDER BY id",
+          s"SELECT id, s.inner.b, s IS NULL FROM $table ORDER BY id",
+          s"SELECT id, items.x FROM $table ORDER BY id",
+          s"SELECT id, m['k'].v FROM $table ORDER BY id",
+          s"SELECT id FROM $table WHERE s.inner.b > 100 ORDER BY id",
+          s"SELECT id, s FROM $table ORDER BY 
id").foreach(checkIcebergNativeScan)
+
+        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. A nested
+  // partition source that the query prunes away makes the task read with the 
full schema.
+  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") {
+
+        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, s STRUCT<a: INT, pad: STRING>) USING 
iceberg $morProperties")
+        spark.sql(s"INSERT INTO $mor SELECT id, s FROM ($rows)")
+        val snapshotBeforeDeletes = spark
+          .sql(s"SELECT snapshot_id FROM $mor.snapshots ORDER BY committed_at 
DESC LIMIT 1")
+          .collect()(0)
+          .getLong(0)
+        spark.sql(s"DELETE FROM $mor WHERE id % 10 = 0")
+        commitEqualityDelete("test_cat", "db", "nested_pruning_mor", "id", 7, 
warehouseDir)
+        checkIcebergNativeScan(s"SELECT id, s.a FROM $mor ORDER BY id")
+        // The equality-delete key `id` is not projected.
+        checkIcebergNativeScan(s"SELECT s.a FROM $mor ORDER BY s.a")
+        checkIcebergNativeScan(
+          s"SELECT id, s.a FROM $mor VERSION AS OF $snapshotBeforeDeletes 
ORDER BY id")

Review Comment:
   This also fixes a case the tests don't pin down. Before, a task with deletes 
read with `task.schema()`, which is the current table schema. So a `VERSION AS 
OF` read of a column dropped after that snapshot failed natively with `field 
not found` whenever the table had deletes. With the pruned scan schema it 
works. I checked `SELECT id, c FROM t VERSION AS OF <snapshot>` with a position 
delete and `c` dropped afterwards. It fails with the config off and matches 
Spark with it on. Could this test cover that case so it stays fixed? Could it 
also assert that these delete tasks get the pruned schema, the way the pool 
test checks for `pad`? Right now the results would still match if every task 
with deletes went back to the full schema.



##########
spark/src/test/scala/org/apache/comet/CometIcebergNativeSuite.scala:
##########
@@ -2284,6 +2285,301 @@ 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. NULL
+  // structs, lists, maps, and list elements check that validity is rebuilt 
from the pruned leaves.
+  // 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, inner: STRUCT<b: 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),
+            IF(id % 7 = 0, NULL, named_struct(
+              'a', IF(id % 5 = 0, NULL, CAST(id AS INT)),
+              'pad', pad,
+              'inner', named_struct('b', CAST(id * 2 AS INT), 'pad', pad))),
+            IF(id % 7 = 0, NULL, array(
+              named_struct('x', CAST(id AS INT), 'pad', pad),
+              IF(id % 5 = 0, NULL, named_struct('x', CAST(-id AS INT), 'pad', 
pad)))),
+            IF(id % 7 = 0, NULL, 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))
+        """)
+
+        Seq(
+          s"SELECT id, s.a FROM $table ORDER BY id",
+          s"SELECT id, s.inner.b, s IS NULL FROM $table ORDER BY id",
+          s"SELECT id, items.x FROM $table ORDER BY id",
+          s"SELECT id, m['k'].v FROM $table ORDER BY id",
+          s"SELECT id FROM $table WHERE s.inner.b > 100 ORDER BY id",
+          s"SELECT id, s FROM $table ORDER BY 
id").foreach(checkIcebergNativeScan)
+
+        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. A nested
+  // partition source that the query prunes away makes the task read with the 
full schema.
+  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") {
+
+        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, s STRUCT<a: INT, pad: STRING>) USING 
iceberg $morProperties")
+        spark.sql(s"INSERT INTO $mor SELECT id, s FROM ($rows)")
+        val snapshotBeforeDeletes = spark
+          .sql(s"SELECT snapshot_id FROM $mor.snapshots ORDER BY committed_at 
DESC LIMIT 1")
+          .collect()(0)
+          .getLong(0)
+        spark.sql(s"DELETE FROM $mor WHERE id % 10 = 0")
+        commitEqualityDelete("test_cat", "db", "nested_pruning_mor", "id", 7, 
warehouseDir)
+        checkIcebergNativeScan(s"SELECT id, s.a FROM $mor ORDER BY id")
+        // The equality-delete key `id` is not projected.
+        checkIcebergNativeScan(s"SELECT s.a FROM $mor ORDER BY s.a")
+        checkIcebergNativeScan(
+          s"SELECT id, s.a FROM $mor VERSION AS OF $snapshotBeforeDeletes 
ORDER BY id")
+
+        val partitioned = "test_cat.db.nested_pruning_partitioned"
+        spark.sql(s"""
+          CREATE TABLE $partitioned (id INT, p STRING, s STRUCT<a: INT, pad: 
STRING>)
+          USING iceberg PARTITIONED BY (p) $morProperties
+        """)
+        spark.sql(s"INSERT INTO $partitioned $rows")
+        spark.sql(s"DELETE FROM $partitioned WHERE id % 10 = 0")
+        // The partition source `p` is not projected.
+        checkIcebergNativeScan(s"SELECT id, s.a FROM $partitioned ORDER BY id")
+        checkIcebergNativeScan(s"SELECT p, count(s.a) FROM $partitioned GROUP 
BY p ORDER BY p")
+
+        // The top-level `region` would collide with `s.region` appended at 
the top level.
+        val nestedSource = "test_cat.db.nested_pruning_nested_source"
+        spark.sql(s"""
+          CREATE TABLE $nestedSource (
+            id INT, region STRING, s STRUCT<region: STRING, a: INT, pad: 
STRING>)
+          USING iceberg PARTITIONED BY (s.region)
+        """)
+        spark.sql(s"""
+          INSERT INTO $nestedSource
+          SELECT CAST(id AS INT), 'top', named_struct('region', IF(id % 2 = 0, 
'east', 'west'),
+            'a', CAST(id AS INT), 'pad', repeat('x', 100))
+          FROM range(200)
+        """)
+        checkIcebergNativeScan(s"SELECT id, region, s.a FROM $nestedSource 
ORDER BY id")
+        checkIcebergNativeScan(s"SELECT id, s.region FROM $nestedSource ORDER 
BY id")

Review Comment:
   I tried this one locally, and I don't think it can be tested against Spark 
yet. With an equality delete keyed on `s.k`, the native scan fails with `field 
not found` for every query on the table, with the config on or off and even 
when the query projects `s.k`. iceberg-rust resolves equality ids against 
top-level fields only. Spark is wrong too when the query prunes `s.k`, because 
`DeleteFilter.fileProjection` only knows how to add a missing key at the top 
level. On Iceberg 1.5.2 it throws `Cannot find required field for ID`, and on 
1.11 the query returns zero rows. So the native failure predates this PR. I 
think the fix is a `CometScanRule` fallback for nested equality keys, and I 
filed #6782 for it rather than hold this PR on it.



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