andygrove commented on code in PR #6578:
URL: https://github.com/apache/datafusion-comet/pull/6578#discussion_r4209696601


##########
spark/src/main/scala/org/apache/spark/sql/comet/execution/arrow/ArrowWriters.scala:
##########
@@ -612,6 +633,37 @@ private[arrow] class StructWriter(
     children: Array[ArrowFieldWriter])
     extends ArrowFieldWriter {
 
+  private val declaredType = 
Utils.fromArrowField(valueVector.getField).asInstanceOf[StructType]
+  private var allowValidatedTrailingFields = false
+
+  override def writeColumnSlice(input: ColumnVector, startRow: Int, numRows: 
Int): Unit = {
+    input.dataType() match {
+      case sourceType: StructType =>
+        require(
+          sourceType.length >= declaredType.length,
+          s"Cannot write ${declaredType.length} fields of struct $name from " +
+            s"${sourceType.length} fields")
+        val sourcePrefix = 
StructType(sourceType.fields.take(declaredType.length))
+        // Compare in Arrow's type domain. Spark logical annotations that 
Arrow cannot represent
+        // (for example string collations) are intentionally erased by the 
same conversion that
+        // produced the declared Arrow field.
+        val normalizedSourcePrefix = Utils
+          .fromArrowField(Utils.toArrowField(name, sourcePrefix, nullable = 
true, "UTC"))
+          .asInstanceOf[StructType]
+        require(
+          DataType.equalsIgnoreCompatibleNullability(normalizedSourcePrefix, 
declaredType),
+          s"Cannot write struct $name with declared type 
${declaredType.simpleString} " +
+            s"from source type ${sourceType.simpleString}: leading fields are 
incompatible")
+      case sourceType =>
+        throw new IllegalArgumentException(
+          s"Cannot write Arrow struct $name from non-struct source type 
${sourceType.simpleString}")
+    }
+
+    allowValidatedTrailingFields = true

Review Comment:
   After the rebase this flag won't mean top level anymore. On main, 
`ArrayWriter`, `MapWriter` and `StructWriter` all call their children's 
`writeColumnSlice`, so nested struct writers get here too. An `array<struct>` 
element with extra trailing fields would be trimmed rather than rejected, and 
`nested array struct writes require an exact nested shape` would still pass 
because it goes through `RowArrowReader`. `ArrayWriter` also calls 
`elementWriter.writeColumnSlice` once per run of adjacent child rows. Spark's 
nested Parquet reader leaves a slot for each null or empty array, which breaks 
the runs, so the Arrow round trip in this check could run nearly once per row. 
Could the permission to trim be a constructor argument that only 
`ArrowWriter.create` sets for top-level columns, with the check done once per 
input vector rather than per slice?



##########
spark/src/main/scala/org/apache/spark/sql/comet/execution/arrow/ArrowWriters.scala:
##########
@@ -612,6 +633,37 @@ private[arrow] class StructWriter(
     children: Array[ArrowFieldWriter])
     extends ArrowFieldWriter {
 
+  private val declaredType = 
Utils.fromArrowField(valueVector.getField).asInstanceOf[StructType]
+  private var allowValidatedTrailingFields = false
+
+  override def writeColumnSlice(input: ColumnVector, startRow: Int, numRows: 
Int): Unit = {
+    input.dataType() match {
+      case sourceType: StructType =>
+        require(
+          sourceType.length >= declaredType.length,
+          s"Cannot write ${declaredType.length} fields of struct $name from " +
+            s"${sourceType.length} fields")
+        val sourcePrefix = 
StructType(sourceType.fields.take(declaredType.length))
+        // Compare in Arrow's type domain. Spark logical annotations that 
Arrow cannot represent
+        // (for example string collations) are intentionally erased by the 
same conversion that
+        // produced the declared Arrow field.
+        val normalizedSourcePrefix = Utils
+          .fromArrowField(Utils.toArrowField(name, sourcePrefix, nullable = 
true, "UTC"))
+          .asInstanceOf[StructType]
+        require(
+          DataType.equalsIgnoreCompatibleNullability(normalizedSourcePrefix, 
declaredType),

Review Comment:
   `equalsIgnoreCompatibleNullability` compares field names at every level, so 
this rejects more than the trimmed case. An equal-width struct whose nested 
names differ from the plan now fails too. Delta's column mapping produces 
exactly that. `DeltaParquetFileFormat.buildReaderWithPartitionValues` gives the 
Parquet reader `prepareSchemaForRead(requiredSchema)`, which renames nested 
fields to their physical names, so the vectors carry `col-...` names while 
`child.schema` has the logical ones. Delta's format is a `ParquetFileFormat`, 
so `spark.comet.convert.parquet.enabled` reaches it, and the cache serializer 
can too. Those batches convert by position on 1.1 and on main, and with this 
change they fail with "leading fields are incompatible". I tried a struct whose 
only field is named `col-5f2a` under a declared `struct<id:int>` at this head 
and it throws. Could the name check apply only where we actually drop trailing 
fields, with `DataType.equalsIgnoreNameAndCompatibleNullability` cov
 ering the rest? It exists from 3.4 through 4.2. A test with a renamed nested 
field would pin it down.



##########
spark/src/main/scala/org/apache/spark/sql/comet/execution/arrow/ArrowWriters.scala:
##########
@@ -44,7 +44,25 @@ import org.apache.spark.unsafe.Platform
 private[arrow] object ArrowWriter {
   def create(root: VectorSchemaRoot, fixedWidthCapacity: Int): ArrowWriter = {
     require(fixedWidthCapacity >= 0, "Fixed-width capacity must be 
non-negative")
-    val children = root.getFieldVectors().asScala.map { vector =>
+    val declaredFields = root.getSchema.getFields.asScala
+    val fieldVectors = root.getFieldVectors().asScala
+    require(
+      fieldVectors.length == declaredFields.length,
+      s"Arrow root has ${fieldVectors.length} vectors for 
${declaredFields.length} fields")
+    val children = fieldVectors.zip(declaredFields).map { case (vector, 
declaredField) =>
+      vector match {
+        case struct: StructVector =>
+          val declaredChildren = declaredField.getChildren
+          if (struct.size() == 0 && !declaredChildren.isEmpty) {
+            struct.initializeChildrenFromFields(declaredChildren)

Review Comment:
   Which caller passes a root whose struct vectors have no children? Every root 
I could find comes from `VectorSchemaRoot.create` or 
`ArrowReader.getVectorSchemaRoot`, and both create vectors with 
`Field.createVector`, which already calls `initializeChildrenFromFields`. The 
test has to build the vector with `createNewSingleVector` to get into this 
state. If it's for the MOR work in #6240, could it move there along with the 
code that builds such a root? As it stands it only covers top-level structs, 
and a struct inside a list would hit the new `require` in `createFieldWriter` 
instead.



##########
spark/src/test/scala/org/apache/spark/sql/comet/execution/arrow/CometArrowStreamSuite.scala:
##########
@@ -1081,6 +1081,255 @@ class CometArrowStreamSuite extends AnyFunSuite with 
Matchers {
     }
   }
 
+  test("columnar struct writes trim compatible trailing fields outside the 
Arrow schema") {

Review Comment:
   After the rebase, this test, `columnar struct validation accepts Spark types 
normalized by Arrow` and `Spark columnar reader repeatedly writes wider 
compatible structs` pass on main without the fix, because #6566's bulk path 
already trims. Could you add the shapes that still fail on main? One is a wider 
struct with a null row and an array or map field, so `writeColumnSlice` takes 
the row path. The other is a wider struct from a `ColumnVector` that isn't 
`OnHeapColumnVector` or `OffHeapColumnVector`, with a null. An `array<struct>` 
column built from Spark vectors would cover the nested case on the columnar 
path.



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