comphead commented on code in PR #6725:
URL: https://github.com/apache/datafusion-comet/pull/6725#discussion_r4223971155
##########
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:
Agreed, the spec allows 2147483447 as a data id. Fixed in 4125e4291a:
`withoutMetadataColumns` now keeps ids up to and including `MaxDataFieldId`.
`CometIcebergNativeWrite` already had that constant with the same meaning, so
it now lives in `IcebergReflection` and both paths use it.
##########
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:
Done in 4125e4291a. `CometIcebergNativeScanMetadata` now carries
`prunedScanSchema`, built in `extract` next to `scanSchema`, so a reflection
failure there reaches `CometScanRule`'s catch and falls back. Serialization
only chooses among the schemas the metadata already holds.
##########
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:
You're right, a failed lookup read as safe to prune. In 4125e4291a
`findFieldObject` lets reflection failures propagate and returns `None` only
when the schema has no such field. All of its callers run during serialization
(`schemaWithRequiredFields`, `isNestedField`, and this guard), so a failure now
fails the query. That also makes `schemaWithRequiredFields` match its doc.
--
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]