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]