viirya commented on code in PR #58050:
URL: https://github.com/apache/spark/pull/58050#discussion_r3798885082


##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetFilters.scala:
##########
@@ -718,6 +962,37 @@ class ParquetFilters(
     // Probably I missed something and obviously this should be changed.
 
     predicate match {
+      // Shredded-variant paths (e.g. "v.`0`"). Only comparison predicates 
that use min/max
+      // statistics are eligible. Each pushes or(leafPredicate, 
isNotNull(residual)...) over every
+      // residual `value` column along the path (see `makeShreddedFilter`). IS 
NULL / IS NOT NULL on
+      // the logical variant field are intentionally out of scope: "the 
extracted field is null" is
+      // not the same as "typed_value is null", so we must not conflate them.
+      case sources.EqualTo(name, value) if canMakeShreddedFilterOn(name, 
value) =>

Review Comment:
   Good catch, confirmed -- thank you. `Not(EqualTo("v.`0`", 700))` recursed 
through the generic `Not` case into the shredded branch, and `not(or(eq(leaf), 
notEq(residual)))` gets rewritten by `LogicalInverseRewriter` into 
`and(notEq(leaf), eq(residual, null))`, which drops the row group whenever the 
residual has no nulls -- the exact unsound AND shape. Reproduced with your `a 
tinyint` / {500, 600} case.
   
   Fixed in a91401f: a negated predicate that references a shredded-variant 
path is no longer pushed. Added `referencesShreddedName` and guard the generic 
`Not` in both `createFilterHelper` and `convertibleFiltersHelper` to return 
`None`; a non-negated shredded conjunct inside an `AND` still pushes. Added a 
unit test (`!=` / `NOT IN` / `Not(Gt)` not pushed, and the AND case) and an 
integration test that runs your `!=` and `NOT IN` repro over an all-fallback 
row group and asserts {500, 600} come back.



##########
sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/ParquetFilters.scala:
##########
@@ -128,6 +138,198 @@ class ParquetFilters(
       fieldNames: Array[String],
       fieldType: ParquetSchemaType)
 
+  /**
+   * Holds the mapping from a logical shredded-variant path (e.g. "v.`0`") to 
the physical
+   * shredded columns needed to push a sound row-group-skipping predicate.
+   *
+   * @param leaf the physical `typed_value` scalar leaf carrying min/max 
statistics
+   * @param residualFieldNames the untyped `value` residual columns along the 
path, from the
+   *                           top-level residual down to the leaf's own-level 
sibling. Each is a
+   *                           physical field-name array. Only residuals that 
exist in this file's
+   *                           schema are included; a value for the path can 
only be hiding in one
+   *                           of these residuals when the typed leaf is NULL, 
so the pushed
+   *                           predicate OR-s in an IS NOT NULL guard on each 
(see
+   *                           `makeShreddedFilter`).
+   */
+  private case class ShreddedVariantField(
+      leaf: ParquetPrimitiveField,
+      residualFieldNames: Seq[Array[String]])
+
+  // Maps logical shredded-variant paths produced by PushVariantIntoScan (e.g. 
"v.`0`") to the
+  // physical shredded columns. Populated only when `variantExtractionSchema` 
is provided and the
+  // physical file schema actually shreds the requested path.
+  //
+  // Soundness: shredding is per-row and per-file best-effort. A row whose 
value does not fit the
+  // shredded type (type mismatch or overflow), or whose field is not shredded 
in this file, is
+  // stored in an untyped `value` residual with `typed_value` NULL. Parquet 
min/max excludes NULLs,
+  // so pushing the predicate on the typed leaf alone could skip a row group 
that still holds a
+  // matching row in a residual. To stay sound we push `or(leafPredicate, 
isNotNull(residual)...)`
+  // over every residual `value` column along the path: Parquet drops the row 
group only when the
+  // leaf cannot match AND every residual is entirely NULL, so a row group is 
skipped only when
+  // every value for the path is provably in the typed leaf. See 
`makeShreddedFilter`.
+  private val nameToShreddedVariantField: Map[String, ShreddedVariantField] = {
+    variantExtractionSchema match {
+      case Some(variantSchema) =>
+        val entries = shreddedVariantEntries(
+          variantSchema.fields.toSeq, schema.asGroupType(), Array.empty, 
Array.empty)
+        if (caseSensitive) {
+          entries.toMap
+        } else {
+          // Mirror `nameToParquetField`: drop names that are ambiguous under 
case-insensitive
+          // matching rather than risk pushing a filter on the wrong physical 
column.
+          val dedup = entries
+            .groupBy(_._1.toLowerCase(Locale.ROOT))
+            .filter(_._2.size == 1)
+            .transform((_, v) => v.head._2)
+          CaseInsensitiveMap(dedup)
+        }
+      case None => Map.empty
+    }
+  }
+
+  // Look up a child of `group` by name, honoring `caseSensitive`. Returns the 
child type together
+  // with its actual physical name so callers build paths from the on-disk 
names (needed for
+  // correct case-insensitive matching, where the requested key case may 
differ from the file's).
+  private def findChild(group: GroupType, name: String): Option[Type] = {
+    group.getFields.asScala.find { f =>
+      if (caseSensitive) f.getName == name else 
f.getName.equalsIgnoreCase(name)
+    }
+  }
+
+  // Look up the untyped `value` residual sibling in `group`, if it exists as 
a non-REPEATED
+  // primitive. Returns the physical field name.
+  private def residualIn(group: GroupType): Option[String] = findChild(group, 
VALUE).collect {
+    case p: PrimitiveType if p.getRepetition != Repetition.REPEATED => 
p.getName
+  }
+
+  // Copy of `getNormalizedLogicalType` from the `nameToParquetField` closure, 
needed here for the
+  // shredded leaf resolution which runs outside that closure.
+  private def getNormalizedLogicalType(p: PrimitiveType): 
LogicalTypeAnnotation = {
+    (p.getPrimitiveTypeName, p.getLogicalTypeAnnotation) match {
+      case (INT32, intType: IntLogicalTypeAnnotation)
+        if intType.getBitWidth() == 32 && intType.isSigned() => null
+      case (INT64, intType: IntLogicalTypeAnnotation)
+        if intType.getBitWidth() == 64 && intType.isSigned() => null
+      case (_, otherType) => otherType
+    }
+  }
+
+  // Navigate the regular shredding layout from a variant column's physical 
group, resolving both
+  // the typed leaf and the residual `value` columns along the path. The 
layout is:
+  //   <col> / typed_value / k0 / typed_value / ... / kN / typed_value   (leaf)
+  //   <col> / value                                                     (L0 
residual)
+  //   <col> / typed_value / k0 / value                                  (L1 
residual)
+  //   ...
+  //   <col> / typed_value / k0 / ... / kN / value                       
(leaf-level residual)
+  // Paths are built from the on-disk field names (via `findChild`) so 
case-insensitive matching
+  // uses the file's actual names. A value for the path can only be hiding in 
one of these residual
+  // `value` columns when the typed leaf is NULL, so IS NULL on all of them is 
the soundness guard.
+  // Residuals absent in this file's schema are skipped (that level cannot 
hold a fallback here).
+  // Returns None if the file does not shred this path down to a non-REPEATED 
scalar leaf (nothing
+  // is pushed and the row group is simply read).
+  private def resolveShredded(
+      physCol: GroupType,
+      physColPath: Array[String],
+      keys: Array[String]): Option[ShreddedVariantField] = {
+    if (keys.isEmpty) return None
+    val residuals = scala.collection.mutable.ArrayBuffer.empty[Array[String]]
+    // L0: the variant column's own residual.
+    residualIn(physCol).foreach(r => residuals += (physColPath :+ r))
+    // Descend key by key: <group>/typed_value/<key>. Collect each level's 
residual sibling.
+    var group = physCol
+    var namePath = physColPath
+    var idx = 0
+    while (idx < keys.length) {
+      val typedChild = findChild(group, TYPED_VALUE) match {
+        case Some(g: GroupType) => g
+        case _ => return None
+      }
+      val typedName = typedChild.getName
+      val keyChild = findChild(typedChild, keys(idx)) match {

Review Comment:
   Agreed, confirmed -- variant keys are data and the reader resolves them 
exact-case (`objectSchemaMap.get` / `getFieldByKey`), so applying 
`caseSensitiveAnalysis` to them can bind the predicate to the wrong subtree 
when a file shreds sibling keys differing only in case.
   
   Fixed in a91401f: `findChild` takes an `exact` flag; object keys and the 
structural `typed_value`/`value` names are matched exact-case, while the 
top-level variant column name still honors `caseSensitive` (it is a Spark 
identifier). Added a unit test that a `$.a` request does not case-insensitively 
bind to a physical `A` subtree under `caseSensitive=false`.



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