[ 
https://issues.apache.org/jira/browse/SPARK-59108?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Yang Jie resolved SPARK-59108.
------------------------------
    Fix Version/s: 4.3.0
       Resolution: Fixed

Issue resolved by pull request 58513
[https://github.com/apache/spark/pull/58513]

> Avro positionalFieldMatching resolves fields against the pruned schema, so a 
> pruned read takes the wrong field
> --------------------------------------------------------------------------------------------------------------
>
>                 Key: SPARK-59108
>                 URL: https://issues.apache.org/jira/browse/SPARK-59108
>             Project: Spark
>          Issue Type: Bug
>          Components: SQL
>    Affects Versions: 4.2.0, 4.3.0, 5.0.0, 4.1.3, 4.0.4
>            Reporter: Yang Jie
>            Assignee: Yang Jie
>            Priority: Major
>              Labels: pull-request-available
>             Fix For: 4.3.0
>
>
> With positionalFieldMatching=true, the Avro reader resolves a catalyst field 
> to an Avro field by position in the *projected* schema rather than in the 
> table's schema, so a column-pruned read takes the wrong Avro field and 
> returns wrong values with no error.
> Reproduction, spark-shell:
> {code}
> val dir = "/tmp/spark-avro-positional"
> spark.range(0, 5).selectExpr("id AS a", "id * 100 AS b", "id * 10000 AS c")
>   .write.mode("overwrite").format("avro").save(dir)
> spark.read.format("avro").option("positionalFieldMatching", "true").load(dir)
>   .createOrReplaceTempView("t")
> sql("SELECT sum(a), sum(b), sum(c) FROM t").show()  // 10, 1000, 100000 -- 
> all correct
> sql("SELECT sum(c) FROM t").show()                  // 10      -- should be 
> 100000
> sql("SELECT sum(b) FROM t").show()                  // 10      -- should be 
> 1000
> sql("SELECT sum(a), sum(c) FROM t").show()          // 10, 1000 -- sum(c) 
> should be 100000
> sql("SELECT sum(a), sum(b) FROM t").show()          // 10, 1000 -- correct
> {code}
> Every value above is measured. Only the read that projects all three fields 
> returns all three correctly: a column's value depends on which other columns 
> the query selects. The reads that happen to be right are the ones whose 
> projection is a prefix of the file's field list, [a] and [a, b].
> When it bites. All of the following have to hold. The read sets 
> positionalFieldMatching (an Avro read option, default false; there is no 
> session configuration for it). The query prunes columns, so the projection is 
> a strict subset of the file's fields. And some projected column sits at a 
> different index in the projection than in the file's field list, which is to 
> say the projection is not a prefix of that list in order. Whether the failure 
> is silent depends on the types of the mispaired fields: matching types return 
> wrong values, as above, and incompatible ones fail the read with a 
> schema-incompatibility error instead.
> Mechanism. AvroPartitionReaderFactory builds the deserializer as 
> AvroDeserializer(userProvidedSchema.getOrElse(reader.getSchema), 
> readDataSchema, positionalFieldMatching, ...). The Avro side is the full 
> schema, never pruned; the catalyst side is readDataSchema, the projection. 
> Under positional matching AvroUtils.AvroSchemaHelper.getAvroField is 
> avroFieldArray.lift(catalystPos), so catalyst field i of the projection pairs 
> with Avro field i of the full schema. Pruning c down to the only projected 
> column makes it catalyst position 0, which pairs with Avro field a. Only 
> validateNoExtraCatalystFields is called on the read side, so nothing rejects 
> the mismatch. AvroFileFormat takes the same path on V1.
> A fix has to pair the catalyst field with its position in the table's data 
> schema rather than in the projection, the way OrcUtils.requestedColumnIds 
> does for orc.force.positional.evolution: it maps through 
> dataSchema.fieldIndex(name) and reports canPruneCols = false, which makes 
> ORC's positional path projection-independent.
> The option itself dates to SPARK-34365 (3.2.0). Everything above is measured 
> on master.
> Context. Found while reviewing SPARK-57205, which lets Spark fuse two DSv2 
> file scans that differ only in their projected columns. Merging does not 
> introduce this bug and cannot turn a correct answer into a wrong one here, 
> because the union of two prefixes is the longer prefix, but it does make the 
> wrong value visible in a new place: two subqueries selecting a and c 
> respectively answer 10 and 10 unmerged and 10 and 1000 merged. SPARK-57205 
> therefore withholds the SCAN_MERGING capability from AvroTable under this 
> option, so that the capability's contract holds for every source that 
> declares it. That gate can be removed once this is fixed.



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to