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


##########
spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala:
##########
@@ -749,6 +749,40 @@ object IcebergReflection extends Logging {
     }
   }
 
+  /** Iceberg reserves field ids from `Integer.MAX_VALUE - 200` up for 
metadata columns. */
+  private val MinReservedFieldId = Int.MaxValue - 200
+
+  /**
+   * Returns `schema` without its top-level metadata columns (`_file`, `_pos`, 
`_partition` and
+   * the other reserved ids), or `schema` itself when it has none. A scan's 
expected schema
+   * includes the metadata columns the query selects, which a table schema 
never has.
+   */
+  def withoutMetadataColumns(schema: Any): Any = {
+    import scala.jdk.CollectionConverters._
+    val columns =
+      getMethod(schema.getClass, 
"columns").invoke(schema).asInstanceOf[java.util.List[_]]
+    val dataColumns = columns.asScala.filter(fieldIdOf(_) < MinReservedFieldId)

Review Comment:
   The spec says Iceberg tables [must not use field ids greater than 
2147483447](https://github.com/apache/iceberg/blob/f74aea1e68fc8161905a748f069c055f47ca64b5/format/spec.md?plain=1#L439-L441)
 (`Integer.MAX_VALUE - 200`), so 2147483447 itself is a valid id for a data 
column. As I read it, `fieldIdOf(_) < MinReservedFieldId` drops a column with 
that id from the task schema. Iceberg Java's `MetadataColumns` 
[comment](https://github.com/apache/iceberg/blob/f74aea1e68fc8161905a748f069c055f47ca64b5/core/src/main/java/org/apache/iceberg/MetadataColumns.java#L63)
 calls `Integer.MAX_VALUE - (101-200)` reserved, which disagrees with the spec 
at that one id. Should the filter keep ids up to and including 
`Integer.MAX_VALUE - 200`, to follow the spec?
   
   ```suggestion
     /** The spec reserves field ids above `Integer.MAX_VALUE - 200` for 
metadata columns. */
     private val MaxDataFieldId = Int.MaxValue - 200
   
     /**
      * Returns `schema` without its top-level metadata columns (`_file`, 
`_pos`, `_partition` and
      * the other reserved ids), or `schema` itself when it has none. A scan's 
expected schema
      * includes the metadata columns the query selects, which a table schema 
never has.
      */
     def withoutMetadataColumns(schema: Any): Any = {
       import scala.jdk.CollectionConverters._
       val columns =
         getMethod(schema.getClass, 
"columns").invoke(schema).asInstanceOf[java.util.List[_]]
       val dataColumns = columns.asScala.filter(fieldIdOf(_) <= MaxDataFieldId)
   ```



##########
spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala:
##########
@@ -1016,6 +1016,37 @@ object CometIcebergNativeScan extends 
CometOperatorSerde[CometBatchScanExec] wit
     val pageIndexUnsupportedColumns =
       IcebergReflection.pageIndexUnsupportedColumns(metadata.tableSchema)
 
+    // iceberg-rust reads every leaf of a projected column that the task 
schema holds. The scan
+    // schema is the read schema after Spark's nested schema pruning, so a 
struct, list, or map
+    // column keeps only the nested fields the query uses. The table schema 
keeps all of them,
+    // which reads and decodes the pruned ones only to drop them later. 
Metadata columns resolve
+    // by field id rather than from the task schema, so they are left out, as 
the table schema
+    // leaves them out.
+    val pruneNestedFields = 
CometConf.COMET_ICEBERG_NESTED_SCHEMA_PRUNING_ENABLED.get()
+    lazy val prunedScanSchema: AnyRef =
+      
IcebergReflection.withoutMetadataColumns(metadata.scanSchema).asInstanceOf[AnyRef]
+    // Partition sources and equality-delete keys that the task schema lacks 
are appended at its
+    // top level. A nested one that the query pruned away would land outside 
its struct, so a task
+    // that needs one reads with the full schema instead.
+    val prunableCache = mutable.HashMap[Seq[Int], Boolean]()
+    def prunesNestedFields(requiredFieldIds: Seq[Int]): Boolean =
+      pruneNestedFields && prunableCache.getOrElseUpdate(
+        requiredFieldIds,
+        requiredFieldIds.forall { id =>
+          IcebergReflection.findFieldObject(prunedScanSchema, id).isDefined ||
+          !IcebergReflection.isNestedField(metadata.tableSchema, id)

Review Comment:
   What happens here if a reflection call fails? `findFieldObject` [catches 
every exception and returns 
`None`](https://github.com/apache/datafusion-comet/blob/3815b1595ae1a13f49c29f3de2775c3fb597416f/spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala#L696-L704),
 and `isNestedField` builds on it. If I'm reading this right, a failed lookup 
reads as "absent from the pruned schema and not nested in the table schema", so 
`prunesNestedFields` returns `true` and the task gets the pruned schema. That 
is the case this guard exists to prevent. This runs during serialization, so I 
think a failure here should fail the query rather than pick a schema.
   
   Every caller of `findFieldObject` runs during serialization: 
`schemaWithRequiredFields`, whose 
[doc](https://github.com/apache/datafusion-comet/blob/3815b1595ae1a13f49c29f3de2775c3fb597416f/spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala#L706-L718)
 already says it throws on any reflection error, and the two new calls here. 
Should `findFieldObject` let the exception propagate, the way 
`nestedFieldsAddedOrRenamed` 
[does](https://github.com/apache/datafusion-comet/blob/3815b1595ae1a13f49c29f3de2775c3fb597416f/spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala#L1302-L1304)?



##########
spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeScan.scala:
##########
@@ -1016,6 +1016,37 @@ object CometIcebergNativeScan extends 
CometOperatorSerde[CometBatchScanExec] wit
     val pageIndexUnsupportedColumns =
       IcebergReflection.pageIndexUnsupportedColumns(metadata.tableSchema)
 
+    // iceberg-rust reads every leaf of a projected column that the task 
schema holds. The scan
+    // schema is the read schema after Spark's nested schema pruning, so a 
struct, list, or map
+    // column keeps only the nested fields the query uses. The table schema 
keeps all of them,
+    // which reads and decodes the pruned ones only to drop them later. 
Metadata columns resolve
+    // by field id rather than from the task schema, so they are left out, as 
the table schema
+    // leaves them out.
+    val pruneNestedFields = 
CometConf.COMET_ICEBERG_NESTED_SCHEMA_PRUNING_ENABLED.get()
+    lazy val prunedScanSchema: AnyRef =
+      
IcebergReflection.withoutMetadataColumns(metadata.scanSchema).asInstanceOf[AnyRef]

Review Comment:
   `withoutMetadataColumns` needs only the scan schema, but it runs here in 
`serializePartitions`, after `CometScanRule` has committed the scan to native 
execution. If its reflection fails (`columns()`, `fieldId()`, or the 
`Schema(List)` constructor), the query fails where it could have fallen back to 
Spark. `CometIcebergNativeScanMetadata.extract` 
[says](https://github.com/apache/datafusion-comet/blob/3815b1595ae1a13f49c29f3de2775c3fb597416f/spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala#L2364-L2368)
 it performs all reflection once during planning and returns `None` to fall 
back, and `CometScanRule` [falls 
back](https://github.com/apache/datafusion-comet/blob/3815b1595ae1a13f49c29f3de2775c3fb597416f/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala#L625-L644)
 when it throws. What do you think about computing the pruned schema in 
`extract` and carrying it on `CometIcebergNativeScanMetadata` next to 
`scanSchema`? Serialization would then only choos
 e between schemas it already holds.



##########
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:
   Should we add a case where the equality-delete key is a nested field inside 
a struct the query prunes away, for example an equality delete keyed on `s.k` 
read with `SELECT s.a`? The spec allows [equality 
delete](https://github.com/apache/iceberg/blob/f74aea1e68fc8161905a748f069c055f47ca64b5/format/spec.md?plain=1#L1418-L1424)
 columns nested in structs, and `CometScanRule` lets a primitive nested key 
through. The nested partition source here covers the partition half of 
`prunesNestedFields`. Nothing covers the equality-delete half, which decides 
the schema iceberg-rust applies the delete against. `commitEqualityDelete` sets 
a top-level field on the delete record, so it would need to build the nested 
record for this.



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