[SPARK-59781][PYTHON] Raise FIELD_STRUCT_LENGTH_MISMATCH instead of truncating rows longer than the schema in createDataFrame - #59038
Conversation
…runcating rows longer than the schema in createDataFrame
In PySpark Classic, createDataFrame from an RDD (schema inferred from
the first row by default) or from a local list with verifySchema=False
silently dropped the trailing values of any row that had more values
than the schema, as long as one field needed a Python-side converter
or conversion. Two zips did it: _create_converter.convert_struct
(zip(obj, converters), any string, struct or null field) and
StructType.toInternal (zip(self.fields, obj, ...), any field whose
needConversion() is true, such as date or timestamp). zip stops at the
shorter side, so
spark.createDataFrame(sc.parallelize([("a", 1), ("b", 2, 3)]))
returned [Row(_1='a', _2=1), Row(_1='b', _2=2)] and a Row with an
extra key in a different position, Row(c=3, a="y", b=2), came back as
Row(a='3', b=None). Whether a user saw an error or a wrong result
depended on the column types: with only long columns the same rows
fail in the JVM with STRUCT_ARRAY_LENGTH_MISMATCH, with
verifySchema=True the verifier raises FIELD_STRUCT_LENGTH_MISMATCH,
and Spark Connect raises AXIS_LENGTH_MISMATCH for every variant.
Check the length of a tuple or list row against the number of fields
at both sites and raise PySparkValueError FIELD_STRUCT_LENGTH_MISMATCH,
the error the type verifier already uses for the same input. The dict
and object branches, which look fields up by name, are unchanged, and
so is the pass-through branch that the JVM already checks. Rows with
fewer values than fields, which previously failed in the JVM at the
first action, now fail with the same PySpark error, up front on the
local path and in the Python worker on the RDD path. 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.
Adds two tests to test_types.py, one on an RDD with an inferred schema
and one on a local list with verifySchema=False; the Spark Connect
parity twin skips the RDD test and asserts AXIS_LENGTH_MISMATCH for
the local list.
|
cc @HyukjinKwon @zhengruifeng, this adds a row length check to Quick repro on master, no config needed: from pyspark.sql import Row
sc = spark.sparkContext
spark.createDataFrame(sc.parallelize([("a", 1), ("b", 2, 3)])).collect()
# [Row(_1='a', _2=1), Row(_1='b', _2=2)] the 3 is gone
spark.createDataFrame(sc.parallelize([Row(a="x", b=1), Row(c=3, a="y", b=2)])).collect()
# [Row(a='x', b=1), Row(a='3', b=None)] wrong values, not just a dropped tail
spark.createDataFrame(sc.parallelize([(1, 1), (2, 2, 3)])).collect()
# STRUCT_ARRAY_LENGTH_MISMATCH same rows, only long columns, so no Python converter: the JVM catches itWith this PR the first two raise |
HyukjinKwon
left a comment
There was a problem hiding this comment.
Thanks, this is a clean fix. Both guards sit exactly where zip could drop data: the convert_fields branch of convert_struct, and the _needSerializeAnyField tuple/list branch of StructType.toInternal. The no-conversion branches stay unchanged, and the JVM already rejects mismatches there. Row also takes the positional tuple branch, so the Row(c=3, a="y", b=2) case is covered. Reusing FIELD_STRUCT_LENGTH_MISMATCH keeps the error consistent with the verifySchema=True path. The local-list test covers both sites: the string column goes through the converter, and "y long, d date" skips the converter and reaches toInternal.
Two non-blocking suggestions:
-
Migration guide entry. Because
StructType.toInternalis shared, jobs that currently finish can start failing. That includes non-Arrow Python UDFs returning a longer struct, andapplyInPandasWithState/transformWithStatestate tuples longer than the declared state schema. They were getting truncated values, but they didn't fail.pyspark_upgrade.rstusually records this kind of change (e.g. the 4.2DATA_SOURCE_RETURN_SCHEMA_MISMATCHentry). A short item under a new "Upgrading from PySpark 4.3 to 4.4" section would help users who hit it. -
A test for the shared
toInternalpaths. The description says non-Arrow UDF struct results and state tuples change too, but the new tests only exercisecreateDataFrame. A small case like auseArrow=FalseUDF with return type"a string, d date"returning a 3-tuple, assertingFIELD_STRUCT_LENGTH_MISMATCH, would pin the behavior for those callers.
…red toInternal check Review follow-up: record the behaviour change in pyspark_upgrade.rst under a new 4.3 to 4.4 section, and add a useArrow=False UDF test that returns a struct longer than its declared return type, which pins FIELD_STRUCT_LENGTH_MISMATCH for the StructType.toInternal callers outside createDataFrame.
Thanks for the review @HyukjinKwon . Both done in 72009ab:
On the heading: master is 5.0.0-SNAPSHOT and branch-4.x is 4.4.0, so I went with 4.3 to 4.4 as you suggested, assuming this gets picked to branch-4.x. If it should stay master only I will rename it. |
What changes were proposed in this pull request?
Two places in
python/pyspark/sql/types.pyturn a tuple or list row into the internal tuple withzip:_create_converter.convert_struct(zip(obj, converters), used when a field needs a Python-side converter, i.e. string, struct or null type) andStructType.toInternal(zip(self.fields, obj, self._needConversion), used when a field needs conversion, e.g. date or timestamp).zipstops at the shorter side, so a row with more values than the schema loses its trailing values without any error. This PR checkslen(obj)against the number of fields at both sites and raisesPySparkValueErrorwith the existing error classFIELD_STRUCT_LENGTH_MISMATCH, which is what the type verifier already raises for the same row whenverifySchema=True. The dict and object branches (matched by field name) and the pass-through branches (already checked by the JVM) are unchanged.Why are the changes needed?
In PySpark Classic,
createDataFramefrom an RDD (the schema is inferred from the first row by default) or from a local list withverifySchema=Falsesilently drops data:Before, on master:
The value 3 is gone in the first and third case, and the second case has wrong values in both columns rather than just a dropped tail. Whether a user gets an error or a wrong result depends on the column types: the same rows with only long columns fail in the JVM with
STRUCT_ARRAY_LENGTH_MISMATCH,verifySchema=Truefails withFIELD_STRUCT_LENGTH_MISMATCH, and Spark Connect fails withAXIS_LENGTH_MISMATCHfor every variant. The RDD path with an explicit schema andverifySchema=False(for example"x string, y long, d date"with a four-value row) only reachesStructType.toInternal, which is why both sites need the check.After:
(raised in the Python worker at the first action for the RDD cases, and directly from
createDataFramefor the local list.)SPARK-11868 quoted this zip in 2015, 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, which is why it was resolved as Cannot Reproduce. The longer-row case still truncated.
Does this PR introduce any user-facing change?
Yes, for PySpark Classic, and only for malformed input. A tuple, list or Row with more values than the schema has fields, which was silently truncated when a field needed a converter or conversion, now fails with
FIELD_STRUCT_LENGTH_MISMATCH. The three sibling paths already fail for this input: the no-converter path (STRUCT_ARRAY_LENGTH_MISMATCHfrom the JVM),verifySchema=True(FIELD_STRUCT_LENGTH_MISMATCH) and Spark Connect (AXIS_LENGTH_MISMATCH). A row with fewer values than fields, which failed in the JVM withSTRUCT_ARRAY_LENGTH_MISMATCHat the first action, now fails withFIELD_STRUCT_LENGTH_MISMATCHinstead, up front for a local list and in the Python worker for an RDD. Rows with the right number of values, dict rows and objects are not affected. BecauseStructType.toInternalis shared, two other paths get the same check when a field of the struct needs conversion: struct results of non-Arrow Python UDFs (useArrow=False, orspark.sql.execution.pythonUDF.arrow.enabled=false; the conf defaults to true since 4.2), where a longer result was truncated and a shorter one failed in the JVM; and the state tuples ofapplyInPandasWithStateandtransformWithState, where a state tuple longer than the user-declared state schema was truncated silently and now fails the query withFIELD_STRUCT_LENGTH_MISMATCH. The migration guide gets an "Upgrading from PySpark 4.3 to 4.4" entry for this.How was this patch tested?
Added two tests to
python/pyspark/sql/tests/test_types.py:test_infer_schema_row_length_mismatch(RDD input with the schema inferred from the first row: a longer tuple, a Row with an extra key, a shorter tuple; the error surfaces at the first action) andtest_create_dataframe_row_length_mismatch_without_verification(local list withverifySchema=False: one longer row per site,"x string, y long"for the converter and"y long, d date"forStructType.toInternal, plus a shorter row; exact error class and parameters).test_parity_types.pyskips the RDD test like the other RDD tests and overrides the local-list test to assertAXIS_LENGTH_MISMATCHfor the same inputs. Without the fix both Classic tests fail (the rows come back truncated, or the shorter rows raise the JVM error instead); with the fix they pass, and reverting either guard alone makes the sub-cases for that site fail again. Ran the fulltest_types,test_parity_types,test_dataframe_creation,test_parity_dataframe_creation,test_udfandtest_pandas_grouped_map_with_statemodules locally, and the snippet above on unmodified and patched trees in Classic and Connect sessions, with the Python worker importing the patched code.Review round:
test_udf_struct_result_length_mismatch(auseArrow=FalseUDF with return typea string, d datereturning a 3-tuple) inBaseUDFTestsMixin, so it runs on Classic, Connect, the six Arrow UDF classes (plain, legacy, non-legacy, Classic and Connect) and the transpile parity class. Red on the unmodifiedtypes.py("PythonException not raised", Classic and Connect), green with the change.Was this patch authored or co-authored using generative AI tooling?
Authored by Li Guo, assisted by Claude Code (Fable 5.1).