Skip to content

[SPARK-59781][PYTHON] Raise FIELD_STRUCT_LENGTH_MISMATCH instead of truncating rows longer than the schema in createDataFrame - #59038

Open
NestDream wants to merge 2 commits into
apache:masterfrom
NestDream:fix-struct-row-length-check
Open

NestDream wants to merge 2 commits into
apache:masterfrom
NestDream:fix-struct-row-length-check

Conversation

@NestDream

@NestDream NestDream commented Sep 25, 2026 •

Copy link
Copy Markdown

What changes were proposed in this pull request?

Two places in python/pyspark/sql/types.py turn a tuple or list row into the internal tuple with zip: _create_converter.convert_struct (zip(obj, converters), used when a field needs a Python-side converter, i.e. string, struct or null type) and StructType.toInternal (zip(self.fields, obj, self._needConversion), used when a field needs conversion, e.g. date or timestamp). zip stops at the shorter side, so a row with more values than the schema loses its trailing values without any error. This PR checks len(obj) against the number of fields at both sites and raises PySparkValueError with the existing error class FIELD_STRUCT_LENGTH_MISMATCH, which is what the type verifier already raises for the same row when verifySchema=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, createDataFrame from an RDD (the schema is inferred from the first row by default) or from a local list with verifySchema=False silently drops data:

from pyspark.sql import SparkSession, Row

spark = SparkSession.builder.master("local[1]").getOrCreate()
sc = spark.sparkContext
print(spark.createDataFrame(sc.parallelize([("a", 1), ("b", 2, 3)])).collect())
print(spark.createDataFrame(sc.parallelize([Row(a="x", b=1), Row(c=3, a="y", b=2)])).collect())
print(spark.createDataFrame([("a", 1), ("b", 2, 3)], "x string, y long", verifySchema=False).collect())

Before, on master:

[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)]

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=True fails with FIELD_STRUCT_LENGTH_MISMATCH, and Spark Connect fails with AXIS_LENGTH_MISMATCH for every variant. The RDD path with an explicit schema and verifySchema=False (for example "x string, y long, d date" with a four-value row) only reaches StructType.toInternal, which is why both sites need the check.

After:

pyspark.errors.exceptions.captured.PythonException: ... pyspark.errors.exceptions.base.PySparkValueError: [FIELD_STRUCT_LENGTH_MISMATCH] Length of object (3) does not match with length of fields (2).

(raised in the Python worker at the first action for the RDD cases, and directly from createDataFrame for 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_MISMATCH from 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 with STRUCT_ARRAY_LENGTH_MISMATCH at the first action, now fails with FIELD_STRUCT_LENGTH_MISMATCH instead, 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. Because StructType.toInternal is 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, or spark.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 of applyInPandasWithState and transformWithState, where a state tuple longer than the user-declared state schema was truncated silently and now fails the query with FIELD_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) and test_create_dataframe_row_length_mismatch_without_verification (local list with verifySchema=False: one longer row per site, "x string, y long" for the converter and "y long, d date" for StructType.toInternal, plus a shorter row; exact error class and parameters). test_parity_types.py skips the RDD test like the other RDD tests and overrides the local-list test to assert AXIS_LENGTH_MISMATCH for 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 full test_types, test_parity_types, test_dataframe_creation, test_parity_dataframe_creation, test_udf and test_pandas_grouped_map_with_state modules 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 (a useArrow=False UDF with return type a string, d date returning a 3-tuple) in BaseUDFTestsMixin, 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 unmodified types.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).

…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.
@NestDream

Copy link
Copy Markdown
Author

cc @HyukjinKwon @zhengruifeng, this adds a row length check to _create_converter.convert_struct and StructType.toInternal: both use zip, so on Classic a row longer than the inferred or declared schema was silently truncated whenever a field needed a converter, while the no-converter path, verifySchema=True and Connect all raise for the same rows. Could you take a look when you have a moment?

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 it

With this PR the first two raise FIELD_STRUCT_LENGTH_MISMATCH, the same error verifySchema=True already gives for these rows.

@HyukjinKwon HyukjinKwon left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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:

  1. Migration guide entry. Because StructType.toInternal is shared, jobs that currently finish can start failing. That includes non-Arrow Python UDFs returning a longer struct, and applyInPandasWithState / transformWithState state tuples longer than the declared state schema. They were getting truncated values, but they didn't fail. pyspark_upgrade.rst usually records this kind of change (e.g. the 4.2 DATA_SOURCE_RETURN_SCHEMA_MISMATCH entry). A short item under a new "Upgrading from PySpark 4.3 to 4.4" section would help users who hit it.

  2. A test for the shared toInternal paths. The description says non-Arrow UDF struct results and state tuples change too, but the new tests only exercise createDataFrame. A small case like a useArrow=False UDF with return type "a string, d date" returning a 3-tuple, asserting FIELD_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.
@NestDream

Copy link
Copy Markdown
Author

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:

  1. Migration guide entry. Because StructType.toInternal is shared, jobs that currently finish can start failing. That includes non-Arrow Python UDFs returning a longer struct, and applyInPandasWithState / transformWithState state tuples longer than the declared state schema. They were getting truncated values, but they didn't fail. pyspark_upgrade.rst usually records this kind of change (e.g. the 4.2 DATA_SOURCE_RETURN_SCHEMA_MISMATCH entry). A short item under a new "Upgrading from PySpark 4.3 to 4.4" section would help users who hit it.
  2. A test for the shared toInternal paths. The description says non-Arrow UDF struct results and state tuples change too, but the new tests only exercise createDataFrame. A small case like a useArrow=False UDF with return type "a string, d date" returning a 3-tuple, asserting FIELD_STRUCT_LENGTH_MISMATCH, would pin the behavior for those callers.

Thanks for the review @HyukjinKwon . Both done in 72009ab:

  1. New "Upgrading from PySpark 4.3 to 4.4" section in pyspark_upgrade.rst with one entry covering the three paths that now raise instead of truncating (createDataFrame, non-Arrow UDF struct results, the state value of applyInPandasWithState / transformWithState).
  2. test_udf_struct_result_length_mismatch in BaseUDFTestsMixin: a useArrow=False UDF with return type "a string, d date" returning a 3-tuple, asserting FIELD_STRUCT_LENGTH_MISMATCH. It runs in UDFTests, the Connect parity class, the six Arrow UDF classes (plain, legacy, non-legacy, Classic and Connect) and the transpile parity class. Without the types.py change it fails with "PythonException not raised" on Classic and on Connect.

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.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants