comphead commented on code in PR #6725:
URL: https://github.com/apache/datafusion-comet/pull/6725#discussion_r4223973367
##########
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:
Thanks both. I left this test out. As Andy found, iceberg-rust matches
equality ids only against the delete file's top-level columns
(`build_field_id_to_arrow_schema_map` in `record_batch_transformer.rs`), so a
nested key fails natively with the config on or off, and Spark is wrong too
when the query prunes the key. #6782 tracks the fallback. Until then, the guard
keeps the full table schema for a task whose nested key the query prunes away,
so this PR never moves such a key to the top level. The nested partition source
case, which goes through the same guard, is now in
`sql-tests/iceberg/nested_schema_pruning.sql`.
##########
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:
Thanks for checking it with the config off. The test covers it in
4125e4291a: `c` is dropped after a snapshot taken right after the position
deletes, and `SELECT id, c, s.a ... VERSION AS OF` reads it. Every query in
that test also checks that no task schema keeps `pad`, so it fails if tasks
with deletes, or tasks that append a partition source, go back to the full
schema. The partitioned query also checks that its tasks share one schema,
which is what the separate pool test checked, so I folded that test in. These
checks run on Iceberg 1.8 and later, like the suite's other `commonData`
checks. With the config off, delete tasks still read with the current table
schema, as before this PR.
--
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]