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

ASF GitHub Bot updated SPARK-59781:
-----------------------------------
    Labels: pull-request-available  (was: )

> PySpark Classic createDataFrame silently drops trailing values of rows longer 
> than the schema when a field needs a Python-side converter
> ----------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: SPARK-59781
>                 URL: https://issues.apache.org/jira/browse/SPARK-59781
>             Project: Spark
>          Issue Type: Bug
>          Components: PySpark
>    Affects Versions: 3.5.8, 4.2.0, 4.3.0, 5.0.0, 4.1.3, 4.0.4, 4.4.0
>            Reporter: Li Guo
>            Priority: Major
>              Labels: pull-request-available
>
> In PySpark Classic, {{createDataFrame}} from an RDD (schema inferred from the 
> first row by default) or from a local list with {{verifySchema=False}} 
> silently drops the trailing values of a row that has more values than the 
> schema, whenever at least one field needs a Python-side converter (string, 
> struct or null type) or conversion ({{needConversion()}}, e.g. date or 
> timestamp). No error, no warning; the extra values are gone and the DataFrame 
> looks fine.
> {code:python}
> from pyspark.sql import SparkSession, Row
> spark = SparkSession.builder.master("local[1]").getOrCreate()
> sc = spark.sparkContext
> # schema inferred from the first row; the second row has one value too many
> print(spark.createDataFrame(sc.parallelize([("a", 1), ("b", 2, 
> 3)])).collect())
> # a Row with an extra key in a different position
> print(spark.createDataFrame(sc.parallelize([Row(a="x", b=1), Row(c=3, a="y", 
> b=2)])).collect())
> # local list, explicit schema, verifier switched off
> print(spark.createDataFrame([("a", 1), ("b", 2, 3)], "x string, y long", 
> verifySchema=False).collect())
> # only long columns: the same rows raise
> print(spark.createDataFrame(sc.parallelize([(1, 1), (2, 2, 3)])).collect())
> {code}
> Output on master (db7a6b6a8d5) with the default settings:
> {noformat}
> [Row(_1='a', _2=1), Row(_1='b', _2=2)]
> [Row(a='x', b=1), Row(a='3', b=None)]
> [Row(x='a', y=1), Row(x='b', y=2)]
> pyspark.errors.exceptions.captured.IllegalArgumentException: 
> [STRUCT_ARRAY_LENGTH_MISMATCH] Input row doesn't have expected number of 
> values required by the schema. 2 fields are required while 3 values are 
> provided. SQLSTATE: 2201E
> {noformat}
> The value 3 disappears in the first case, the second case gives wrong values 
> ('3' in column a, None in column b) rather than just a dropped tail, and the 
> third case drops the value as well. The fourth case shows that the same input 
> is rejected as soon as no column needs a converter, so whether a user gets an 
> error or a wrong result depends on the column types. {{verifySchema=True}} 
> raises {{FIELD_STRUCT_LENGTH_MISMATCH}} for the same rows, and Spark Connect 
> raises {{AXIS_LENGTH_MISMATCH}} for every variant.
> Cause: two places build the internal tuple with {{zip}}, which stops at the 
> shorter side. {{_create_converter.convert_struct}} in 
> python/pyspark/sql/types.py does {{tuple(conv(v) for v, conv in zip(obj, 
> converters))}} when any field needs a converter, and 
> {{StructType.toInternal}} does {{tuple(f.toInternal(v) if c else v for f, v, 
> c in zip(self.fields, obj, self._needConversion))}} when any field needs 
> conversion. The RDD path with an explicit schema and {{verifySchema=False}} 
> reaches only the second one, for example {{createDataFrame(rdd, "x string, y 
> long, d date", verifySchema=False)}} with a four-value row. The no-converter 
> branches pass the tuple through unchanged and the JVM ({{EvaluatePython}}) 
> catches the length there.
> SPARK-11868 (2015, resolved Cannot Reproduce in 2016) quoted the same zip, 
> but its example is a Row with a missing key, i.e. a row shorter than the 
> schema; the 2016 comment on that ticket already shows the JVM rejecting it 
> with the length error. The longer-row case reported here still silently 
> truncates.
> Fix: check {{len(obj)}} against the number of fields at both sites and raise 
> {{PySparkValueError}} with the existing error class 
> {{FIELD_STRUCT_LENGTH_MISMATCH}}, the same error the type verifier raises for 
> this input when {{verifySchema=True}}. Rows with fewer values than fields 
> then get the same error instead of the JVM's {{STRUCT_ARRAY_LENGTH_MISMATCH}} 
> at the first action, and the other callers of {{StructType.toInternal}} get 
> the same check: struct results of non-Arrow Python UDFs, and the state tuples 
> of {{applyInPandasWithState}} and {{transformWithState}}, where a tuple 
> longer than the state schema was truncated silently as well. Dict and object 
> rows, which are matched by field name, and the pass-through branches are 
> unchanged.



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