Yang Jie created SPARK-59107:
--------------------------------

             Summary: Subplan merging must not widen the columns a 
projection-sensitive parser reads
                 Key: SPARK-59107
                 URL: https://issues.apache.org/jira/browse/SPARK-59107
             Project: Spark
          Issue Type: Bug
          Components: SQL
    Affects Versions: 5.0.0
            Reporter: Yang Jie


On the DataSource V1 file source path, subplan merging can widen the set of 
columns a scan reads, and for CSV and JSON that changes which records the 
parser treats as malformed. One subquery's result then depends on what a 
sibling subquery projects.

Reproduction, spark-shell on default configurations, so every read below is V1. 
Each shape is run twice, once as-is and once with subplan merging switched off, 
and every value is measured.

{code}
import java.nio.file.{Files, Paths}
def merging(on: Boolean): Unit = 
spark.conf.set("spark.sql.optimizer.excludedRules",
  if (on) "" else "org.apache.spark.sql.execution.planmerging.MergeSubplans")

val csv = "/tmp/spark-v1-merge-csv"
Files.createDirectories(Paths.get(csv))
Files.write(Paths.get(csv, "data.csv"), "0,0\n1,10\n2,BAD\n3,30\n4,40".getBytes)

// 1. DROPMALFORMED: the record malformed in b is dropped for sum(a) too.
spark.read.schema("a long, b long").option("mode", "DROPMALFORMED").csv(csv)
  .createOrReplaceTempView("t1")
val q1 = "SELECT (SELECT sum(a) FROM t1), (SELECT sum(b) FROM t1)"
merging(true);  sql(q1).show()   // [8, 80]
merging(false); sql(q1).show()   // [10, 80]

// 2. PERMISSIVE, the default mode, with the corrupt-record column in the 
schema.
spark.read.schema("a long, b long, _corrupt_record string").option("mode", 
"PERMISSIVE")
  .option("columnNameOfCorruptRecord", "_corrupt_record").csv(csv)
  .createOrReplaceTempView("t2")
val q2 = "SELECT (SELECT count(_corrupt_record) FROM t2 WHERE a >= 0), " +
  "(SELECT sum(b) FROM t2 WHERE a >= 0)"
merging(true);  sql(q2).show()   // [1, 80]
merging(false); sql(q2).show()   // [0, 80]

// 3. FAILFAST with a row carrying fewer tokens than the schema has columns.
val short = "/tmp/spark-v1-merge-short"
Files.createDirectories(Paths.get(short))
Files.write(Paths.get(short, "data.csv"), "0,0\n1,10\n2\n3,30\n4,40".getBytes)
spark.read.schema("a long, b long").option("mode", "FAILFAST").csv(short)
  .createOrReplaceTempView("t3")
val q3 = "SELECT (SELECT sum(a) FROM t3), (SELECT sum(b) FROM t3)"
merging(true);  sql(q3).show()   // throws MALFORMED_RECORD_IN_PARSING
merging(false); sql(q3).show()   // [10, 80]

// 4. ignoreCorruptFiles, parquet: not about parsing, so this reaches every 
format.
val pq = "/tmp/spark-v1-merge-parquet"
spark.range(0, 10).selectExpr("id AS a", "cast(id AS string) AS b")
  .write.mode("overwrite").parquet(pq)
spark.conf.set("spark.sql.files.ignoreCorruptFiles", "true")
spark.read.schema("a long, b long").parquet(pq).createOrReplaceTempView("t4")
val q4 = "SELECT (SELECT sum(a) FROM t4), (SELECT count(b) FROM t4)"
merging(true);  sql(q4).show()   // [null, 0]
merging(false); sql(q4).show()   // [45, 0]
{code}

Shapes 1 to 3 are CSV, and JSON behaves the same way with equivalent records. 
In shape 1 the scan that reads only a never parses b, so the record malformed 
in b is not dropped for sum(a) until the merge widens the read set. In shape 2 
the merged scan parses b, so the corrupt-record column is populated for a row 
the first subquery had counted as clean. In shape 3, with column pruning on, 
UnivocityParser.parsedSchema is the pruned schema, so a one-column scan matches 
a one-token row and only the merged two-column scan trips the tokens.length != 
parsedSchema.length branch. In shape 4 sum(a) touches only healthy data and is 
correct at 45 until the merge widens the read set; FilePartitionReader.next and 
FileScanRDD both route the failure through 
DataSourceUtils.shouldIgnoreCorruptFileException, which accepts any 
RuntimeException, so the rest of that file's rows go with it.

Mechanism. There is no scan-merge operation on the V1 path. Two 
LogicalRelations over the same paths are canonically equal regardless of 
projection, because top-level column pruning happens in physical planning 
(FileSourceStrategy computes readDataColumns from filterAttributes ++ 
projects), so PlanMerger's identical-plan path reuses one of them and the union 
of the columns is produced one level up in mergeNamedExpressions. The 
planmerging package has no per-source predicate on the V1 path; the 
SCAN_MERGING capability gate in PlanMerger applies only to DSv2 scans.

Scope. Declining the merge needs two things. First, a signal reachable from 
HadoopFsRelation.fileFormat saying the format's parsing depends on the 
projected columns; DataSourceUtils.supportNestedPredicatePushdown and 
SchemaPruning.canPruneDataSchema are the closest existing shapes. Second, the 
decline point has to see whether the projections actually differ: the most 
local seam, a (LogicalRelation, LogicalRelation) arm in 
PlanMerger.tryMergePlans, cannot, because only MergeContext threads downward 
and both leaves carry the full schema. Without that distinction a decline also 
gives up the identical-projection case, which is zero-IO pure reuse. Note also 
that "the projections differ" is not the same as "the scan reads more columns", 
since symmetric filter propagation can add columns independently of the 
projections.

The alternative fix, making malformed-record detection independent of the 
projection, changes user-visible behavior for a plain df.select("a") on the 
same file in three of four configurations and would have to be a config-gated 
mode; CSVSuite pins the current behavior for SPARK-29101 (count() is 4 with CSV 
column pruning on and 3 with it off). Turning 
spark.sql.csv.parser.columnPruning.enabled off is not a fix either: it 
addresses only the arity check, while the type-error path is decided by a 
conversion loop over the required schema.

Context. SPARK-57205 declares SCAN_MERGING per format on the DSv2 read path and 
deliberately leaves CSV and JSON out, and withholds it from any file table 
whose reads are not strict, so V2 does not copy any of these behaviours. That 
leaves V1 and V2 disagreeing on these shapes until 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