sunchao commented on code in PR #5365:
URL: https://github.com/apache/datafusion-comet/pull/5365#discussion_r3835076676


##########
contrib/delta-spark/src/main/scala/org/apache/spark/sql/comet/CometDeltaNativeScanExec.scala:
##########
@@ -0,0 +1,250 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.spark.sql.comet
+
+import org.apache.spark.rdd.RDD
+import org.apache.spark.sql.catalyst.expressions._
+import org.apache.spark.sql.catalyst.plans.QueryPlan
+import org.apache.spark.sql.catalyst.plans.physical.{Partitioning, 
UnknownPartitioning}
+import org.apache.spark.sql.execution.{FileSourceScanExec, InSubqueryExec, 
SparkPlan}
+import org.apache.spark.sql.execution.datasources.HadoopFsRelation
+import org.apache.spark.sql.execution.metric.SQLMetric
+import org.apache.spark.sql.types.StructType
+import org.apache.spark.sql.vectorized.ColumnarBatch
+
+import org.apache.comet.contrib.delta.DeltaSparkScanEnvelope
+import org.apache.comet.serde.OperatorOuterClass
+import org.apache.comet.serde.OperatorOuterClass.Operator
+
+/**
+ * Native scan node for Delta Lake tables (contrib). Delta's own planning (log 
replay, snapshot
+ * resolution, partition pruning) has already run inside delta-spark by the 
time this node is
+ * created from the DSv1 [[FileSourceScanExec]]; file listing and split 
planning are delegated to
+ * a [[CometScanExec]] helper, and data reads execute through Comet's native 
DataFusion parquet
+ * machinery, inheriting row-group and page-index pruning.
+ *
+ * DPP: `runtimeFilters` is a constructor field included in equality, so
+ * `CometPlanAdaptiveDynamicPruningFilters`'s rewrite (via 
[[CometScanWithPlanData]]) survives
+ * plan copies, the lesson from CometIcebergNativeScanExec (a transient field 
is dropped by
+ * `TreeNode.makeCopy` on MERGE re-planning).
+ */
+case class CometDeltaNativeScanExec(
+    override val nativeOp: Operator,
+    override val output: Seq[Attribute],
+    requiredSchema: StructType,
+    runtimeFilters: Seq[Expression],
+    dataFilters: Seq[Expression],
+    @transient relation: HadoopFsRelation,
+    originalPlan: FileSourceScanExec,
+    override val serializedPlanOpt: SerializedPlan,
+    sourceKey: String)
+    extends CometLeafExec
+    with CometScanWithPlanData {
+
+  override val nodeName: String = s"CometDeltaNativeScan $relation"
+
+  // Derived from (originalPlan, runtimeFilters), never stored: any copy of 
this node, our
+  // own withDynamicPruningFilters, or a generic Catalyst expression rewrite 
going through
+  // TreeNode.makeCopy, automatically gets a helper consistent with ITS 
runtimeFilters. A
+  // stored helper field would desync from rewritten filters (the #3510 class 
of bug). The
+  // cost is that file listing runs once per executed instance (planning 
listed separately in
+  // the rule extension); correctness over the duplicate driver-side listing.
+  @transient private lazy val scanHelper: CometScanExec =
+    CometDeltaNativeScanExec.planningHelper(originalPlan, runtimeFilters)
+
+  override lazy val outputPartitioning: Partitioning =
+    UnknownPartitioning(perPartitionData.length)

Review Comment:
   **[P2] Avoid executing adaptive pruning while inspecting partitioning**
   
   This getter forces `perPartitionData`, which calls 
`InSubqueryExec.updateResult()`. During AQE, that subquery can still be a 
non-executable adaptive broadcast placeholder.
   
   A reduced Spark 4.0.2 / Delta 4.0.0 planning harness reproduced this through 
Spark's normal AQE validation: a DPP join in one `UNION ALL` branch and a 
coalescible shuffle in another caused validation to inspect this partitioning 
before the custom DPP rewrite. It then failed with 
`CometSubqueryAdaptiveBroadcastExec ... does not support the execute() code 
path`. Other operators remained on Spark, and no native Comet reader executed.
   
   Could we return `UnknownPartitioning(0)` while adaptive placeholders remain 
and make this a non-lazy `def`, so the temporary value is not cached? A 
regression with a query-time dimension filter would help. The current DPP test 
filters the dimension before writing it, so it does not require dynamic pruning.
   



##########
contrib/delta-spark/src/main/scala/org/apache/comet/contrib/delta/CometDeltaNativeScan.scala:
##########
@@ -0,0 +1,533 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.comet.contrib.delta
+
+import scala.jdk.CollectionConverters._
+
+import org.apache.hadoop.fs.Path
+import org.apache.spark.internal.Logging
+import org.apache.spark.sql.catalyst.expressions.Literal
+import org.apache.spark.sql.comet.{CometScanExec, DeltaPlanDataInjector}
+import org.apache.spark.sql.delta.DeltaParquetFileFormat
+import org.apache.spark.sql.delta.RowIndexFilterType
+import org.apache.spark.sql.delta.actions.DeletionVectorDescriptor
+import org.apache.spark.sql.execution.{FileSourceScanExec, ScalarSubquery => 
ExecScalarSubquery}
+import org.apache.spark.sql.execution.datasources.{FilePartition, 
PartitionedFile}
+import org.apache.spark.sql.types.{ByteType, LongType, MetadataBuilder, 
StructField, StructType}
+
+import org.apache.comet.objectstore.NativeConfig
+import org.apache.comet.serde.OperatorOuterClass
+import org.apache.comet.serde.OperatorOuterClass.Operator
+import org.apache.comet.serde.QueryPlanSerde.{exprToProto, serializeDataType}
+import org.apache.comet.serde.operator.{literalToProto, partition2Proto, 
schema2Proto, CometNativeScan}
+import org.apache.comet.shims.ShimFileFormat
+
+/**
+ * Serde for the native Delta scan. Two shapes:
+ *   - Plain reads reuse core's `NativeScanCommon` builder wholesale.
+ *   - Deletion-vector reads: Delta's planner appends 
`__delta_internal_is_row_deleted` (tinyint)
+ *     and Spark's row-index temp column (bigint) to the read schema and 
filters on is_row_deleted
+ *     above the scan. The native reader applies the DV as a row selection, so 
surviving rows are
+ *     by construction not deleted: both internal columns are emitted as 
per-file constants (0),
+ *     the parquet read schema is stripped to the real data columns, and the 
DV descriptor ships
+ *     per file for the native side to fetch and decode.
+ */
+object CometDeltaNativeScan
+    extends Logging
+    with org.apache.spark.sql.catalyst.expressions.PredicateHelper {
+
+  val IsRowDeletedColumn: String = 
DeltaParquetFileFormat.IS_ROW_DELETED_COLUMN_NAME
+  val RowIndexColumn: String = ShimFileFormat.ROW_INDEX_TEMPORARY_COLUMN_NAME
+
+  private[delta] val internalColumnNames: Set[String] = 
Set(IsRowDeletedColumn, RowIndexColumn)
+
+  // Prefix for the internal columns' slots in the partition schema, mirroring 
core's
+  // _comet_metadata_ prefix rationale: DataFusion matches partition columns 
by name.
+  private val deltaConstFieldPrefix = "_comet_delta_"
+
+  def isDvShape(scanExec: FileSourceScanExec): Boolean =
+    scanExec.requiredSchema.exists(f => internalColumnNames.contains(f.name))
+
+  private def deltaFormat(scanExec: FileSourceScanExec): 
DeltaParquetFileFormat =
+    scanExec.relation.fileFormat.asInstanceOf[DeltaParquetFileFormat]
+
+  private def columnMappingMode(scanExec: FileSourceScanExec): String =
+    deltaFormat(scanExec).metadata.columnMappingMode.name
+
+  /**
+   * Under column mapping, parquet files store physical column names (stable 
UUIDs / ids), so the
+   * schemas passed to the native parquet reader must be physical. Positions 
and structure are
+   * preserved, so all positional output binding and projection are 
unaffected. The scan's
+   * internal DV columns are not part of the table schema and must be stripped 
before calling
+   * this.
+   */
+  private def toPhysical(scanExec: FileSourceScanExec, schema: StructType): 
StructType = {
+    val format = deltaFormat(scanExec)
+    if (format.metadata.columnMappingMode.name == "none") {
+      schema
+    } else {
+      // Name mode matches file columns by physical NAME. createPhysicalSchema 
also stamps
+      // parquet.field.id metadata, but files written before the 
column-mapping upgrade have
+      // no field ids and would fail the reader's id expectations, strip the 
ids so the
+      // reader stays purely name-based (id mode, when enabled, will keep 
them).
+      stripFieldIds(org.apache.spark.sql.delta.DeltaColumnMapping
+        .createPhysicalSchema(schema, format.metadata.schema, 
format.metadata.columnMappingMode))
+    }
+  }
+
+  private def stripFieldIds(schema: StructType): StructType = {
+    import org.apache.spark.sql.types._
+    def stripType(dt: DataType): DataType = dt match {
+      case s: StructType => stripFieldIds(s)
+      case a: ArrayType => a.copy(elementType = stripType(a.elementType))
+      case m: MapType =>
+        m.copy(keyType = stripType(m.keyType), valueType = 
stripType(m.valueType))
+      case other => other
+    }
+    StructType(schema.fields.map { f =>
+      val metadata = new MetadataBuilder()
+        .withMetadata(f.metadata)
+        .remove("parquet.field.id")
+        // Sibling key Delta stamps on array/map fields under 
IcebergCompat/Uniform.
+        .remove("parquet.field.nested.ids")
+        .build()
+      f.copy(dataType = stripType(f.dataType), metadata = metadata)
+    })
+  }
+
+  /**
+   * Build the planning-time `DeltaScan` operator (common data only; file 
partitions are injected
+   * lazily at execution). Returns None when an output data type cannot be 
serialized or the plan
+   * shape is not one we can translate faithfully.
+   */
+  def convert(scanExec: FileSourceScanExec, scanHelper: CometScanExec): 
Option[Operator] = {
+    val relation = scanExec.relation
+
+    val firstFileUri = scanHelper.selectedPartitions
+      .flatMap(_.files.headOption)
+      .headOption
+      .map(_.getPath.toUri)
+
+    val hadoopConf = relation.sparkSession.sessionState
+      .newHadoopConfWithOptions(relation.options)
+
+    val tableRootPath = relation.location.rootPaths.head
+    val tableRoot = tableRootPath.toString
+
+    val commonOpt = if (!isDvShape(scanExec)) {
+      // Under column mapping (name mode) the parquet reader must see physical 
names;
+      // positions are preserved so output binding and projection stay 
untouched.
+      CometNativeScan.buildNativeScanCommon(
+        source = scanExec.simpleStringWithNodeId(),
+        output = scanExec.output,
+        requiredSchema = toPhysical(scanExec, scanExec.requiredSchema),
+        dataSchema = toPhysical(scanExec, relation.dataSchema),
+        partitionSchema = relation.partitionSchema,
+        fileConstantMetadataColumns = scanExec.fileConstantMetadataColumns,
+        dataFilters = scanHelper.supportedDataFilters,
+        firstFileUri = firstFileUri,
+        hadoopConf = hadoopConf,
+        conf = scanExec.conf)
+    } else {
+      buildDvScanCommon(scanExec, scanHelper, firstFileUri, hadoopConf)
+    }
+
+    commonOpt.map { commonBuilder =>
+      // Union object-store options over every authority a partition of this 
scan may need a
+      // store for, not just the first data file's scheme (finding 8).
+      val dvDescriptors = DeltaScanSupport.selectedDvDescriptors(scanHelper, 
tableRoot)
+      commonBuilder.putAllObjectStoreOptions(
+        mergedObjectStoreOptions(
+          hadoopConf,
+          storeUris(dvDescriptors, tableRootPath, firstFileUri)).asJava)

Review Comment:
   **[P2] Fall back when native S3 cannot preserve Hadoop credential aliases**
   
   Could we check authentication compatibility before claiming these scans? 
With `SimpleAWSCredentialsProvider` explicitly selected and the access/secret 
keys stored only in JCEKS, Hadoop resolves the credentials through 
`Configuration.getPassword`, but this extraction forwards neither the 
credentials nor the global credential-provider path. The native Simple provider 
therefore has no credentials.
   
   A local probe verified that Hadoop resolves both aliases while the current 
extraction omits them. The underlying native limitation predates this PR, but 
claiming Delta scans exposes it to reads that previously used Hadoop 
successfully.
   
   A conservative admission fallback, with a focused credential-alias test, 
would be sufficient here. Full credential-provider support can follow 
separately.
   



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