comphead commented on code in PR #6725:
URL: https://github.com/apache/datafusion-comet/pull/6725#discussion_r4238758740
##########
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:
You're right, thanks for the repro. Iceberg's `Schema` constructor indexes
the names and rejects the second `c`, and the config-off path goes through the
same `schemaWithRequiredFields`, so the switch doesn't help. In d88698c442 the
missing fields all come from one schema: the current one, or the newest in the
history that has them all under names the task schema does not use. That covers
your drop-and-rename and swap cases, and also the two-source case sunchao
found, where picking a name per field still clashed. It keeps current names and
types when they don't clash, and needs no extra reflection to find the scan's
snapshot. The time travel test covers the drop-and-rename case on the
partitioned table: after `c` is dropped and `p` is renamed to `c`, `SELECT id,
c, s.a ... VERSION AS OF` reads the snapshot's `c`, and `p` is appended under
its older name.
##########
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:
Thanks for finding this and filing #6822. The test covers it since
335e7e0caf: after `snapshotWithDeletes`, `r` is renamed to `r2` and `x` and `y`
swap names, next to the dropped `c`, and `SELECT id, c, r, x, y, s.a ...
VERSION AS OF` reads them under the snapshot's names, with position deletes.
##########
spark/src/main/scala/org/apache/comet/iceberg/IcebergReflection.scala:
##########
@@ -741,28 +743,68 @@ 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)
- val resolvedFields = missingIds.map { id =>
- history.iterator
- .flatMap(s => findFieldObject(s, id))
- .toSeq
- .headOption
- .getOrElse(throw new IllegalStateException(
- s"Cannot resolve field id $id in table schema history"))
- }
val existing =
getMethod(baseSchema.getClass, "columns")
.invoke(baseSchema)
.asInstanceOf[java.util.List[_]]
val newColumns = new java.util.ArrayList[Any](existing)
- resolvedFields.foreach(newColumns.add)
+ // A field is appended under a name it has had that the task schema does
not use yet. The
+ // current schema comes first, to keep a live column's current name and
type, and then
+ // table.schemas(), oldest first. A VERSION AS OF scan schema carries
the snapshot's names,
+ // so a current name can clash with a column the snapshot still has.
+ val schemas = getMethod(table.getClass, "schema").invoke(table) +:
getAllSchemas(table)
+ val names =
scala.collection.mutable.Set(existing.asScala.map(fieldName).toSeq: _*)
+ missingIds.foreach { id =>
+ val (field, name) = schemas.iterator
+ .flatMap(findFieldObject(_, id))
+ .map(f => (f, fieldName(f)))
+ .find { case (_, n) => !names.contains(n) }
Review Comment:
Thanks, this holds from the code. The partition source ids come back in spec
order, `[2, 3]`, id 2 takes its current name `q`, and id 3's names `c` and `q`
are then both taken. Fixed in d88698c442: `schemaWithRequiredFields` now takes
all the missing ids from one schema, the current one or the newest in the
history that has them all under names the task schema does not use. One schema
cannot give two fields the same name, so the order of the ids no longer
matters, and newest first keeps a promoted column's widest type. I searched the
history rather than looking up the scan's snapshot schema, because it needs no
further reflection into Iceberg's scan internals, and whenever the snapshot's
schema has all the ids it qualifies too. The time travel test now has your
case, with both sources bucketed and `c` read through `VERSION AS OF`.
--
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]