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

Yang Jie resolved SPARK-57205.
------------------------------
    Fix Version/s: 4.4.0
       Resolution: Fixed

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

> Declare SCAN_MERGING on the built-in file sources
> -------------------------------------------------
>
>                 Key: SPARK-57205
>                 URL: https://issues.apache.org/jira/browse/SPARK-57205
>             Project: Spark
>          Issue Type: Improvement
>          Components: SQL
>    Affects Versions: 4.3.0, 5.0.0
>            Reporter: Yang Jie
>            Assignee: Yang Jie
>            Priority: Major
>              Labels: pull-request-available
>             Fix For: 4.4.0
>
>
> SPARK-40259 (#57360) made DSv2 scan merging a Spark-side mechanism: 
> PlanMerger drives V2ScanRelationPushDown.rebuildScan to rebuild the merged 
> scan, and a source only declares that widening the set of columns pruned on 
> its builder, with the scan options and pushed filters held constant, changes 
> neither which rows the scan returns nor the values it returns for the columns 
> it was already asked for. It may at most surface a read error. That 
> declaration is the SCAN_MERGING table capability, and no built-in source made 
> it.
> This adds a supportsScanMerging seam to FileTable and overrides it in 
> ParquetTable, OrcTable, TextTable and AvroTable, so those four take part in 
> scan merging on their DSv2 read path. Today two scans of the same file table 
> that differ only in their projected columns are not canonically identical, 
> because a file source folds its data filters into the FileScan object and 
> FileScan.equals compares them. On the V1 path the same filters sit above an 
> identical LogicalRelation, so V1 merges shapes that V2 does not.
> CSVTable and JsonTable deliberately do not override it. Their parsers are 
> handed the columns the scan asked for and decide from that set what counts as 
> a malformed record, so a merged scan reading the union of two column sets can 
> drop or rewrite rows the narrower scan returned. Measured on both read paths: 
> DROPMALFORMED drops a record malformed only in the other subquery's column, 
> PERMISSIVE with a corrupt-record column populates it for a row the narrow 
> scan counted as clean, and FAILFAST with a short CSV row throws where the 
> unmerged query returned rows.
> Two further gates keep the contract true of the four that do declare it. 
> FileTable withholds the capability when a read is not strict, because 
> ignoreCorruptFiles swallows a failure in a column only the sibling projects 
> and drops the rest of that file's rows with it. AvroTable withholds it under 
> positionalFieldMatching, which resolves a column by its position in the 
> projection: SPARK-59108.
> Out of scope, both already listed as follow-ups on SPARK-40259: merging two 
> scans whose strict (fully enforced) filters differ, which is what keeps two 
> differing partition filters apart on the V2 path; and widening nested columns 
> to the union of the fields each side reads, which is what keeps Parquet and 
> ORC nested columns apart while nested schema pruning is on. Two more gaps are 
> recorded on the pull request: differing data filters need 
> dsv2SymmetricFilterPropagation, and CSV and JSON now merge less than V1 does, 
> which is a V1 bug to fix rather than a target to copy and has its own ticket, 
> SPARK-59107.



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