L. C. Hsieh created SPARK-59117:
-----------------------------------
Summary: Set __module__ on the Rust drop-in's PyO3 classes to the
reference pyspark path
Key: SPARK-59117
URL: https://issues.apache.org/jira/browse/SPARK-59117
Project: Spark
Issue Type: Bug
Components: Connect
Affects Versions: connect-rust-4.2.0
Reporter: L. C. Hsieh
Assignee: L. C. Hsieh
The PyO3 classes exposed by the Rust drop-in (pyspark._pyspark) all report
__module__ == "builtins", because that is PyO3's default when #[pyclass]
carries no
`module` argument. None of the 99 #[pyclass] attributes in
crates/pyspark-rs/src set
it. The reference client reports the module that defines each class
(pyspark.sql.types.IntegerType, etc.), so this is a behavioral divergence from
PySpark, not just a cosmetic one.
Two consequences:
1. The class ITSELF is unpicklable. pickle serializes a class BY REFERENCE —
it
writes __module__ + __qualname__ and re-imports on load — so the lookup
lands in
`builtins` and fails:
>>> import pickle
>>> from pyspark.sql.types import IntegerType
>>> pickle.dumps(IntegerType)
PicklingError: Can't pickle <class 'builtins.IntegerType'>:
attribute lookup IntegerType on builtins failed
Real pyspark pickles these classes fine. Reproduced for DataFrame, Column,
and
every DataType class.
2. Error messages, reprs, and any __module__-based dispatch show
"builtins.DataFrame" instead of the pyspark path. That is user-visible and
undercuts the drop-in premise.
Pickling of type INSTANCES already works, but only because each DataType
hand-writes
a __reduce__ that routes through pyspark.sql.types._parse_datatype_json_string
(crates/pyspark-rs/src/types.rs). In other words the cost of the wrong
__module__ is
already being paid — worked around per class rather than fixed at the root.
Those
__reduce__ impls must be KEPT regardless: they carry the type's parameters
(DecimalType(12,3), ArrayType(...)), which a by-reference pickle of the bare
class
would not restore.
Fix: add an explicit `module = "..."` to every #[pyclass] (including the
abstract_type! macro that generates the 8 intermediate type bases), and set
the
module on create_exception! for SkipRestOfInputTableException.
The class-to-module mapping is derived from where the drop-in skin
(python/pyspark/**) already re-exports each class, cross-checked against
Apache Spark
v4.2.0 (commit 32f72996011 — the same commit pinned in
crates/spark-connect-proto/proto/PROTO_SHA.txt).
Two mapping notes worth recording:
* DataFrame and Column map to pyspark.sql.connect.*, NOT pyspark.sql.*.
Upstream's
pyspark/sql/dataframe.py:78 and pyspark/sql/column.py:41 hold the ABSTRACT
parent; the concrete class a Connect session hands the user is
pyspark/sql/connect/dataframe.py:112 and .../connect/column.py:109, which
is what
type(df).__module__ reports on a real Connect session.
* RustRpcError keeps __module__ == "_pyspark": it is an internal transport
error
with no upstream counterpart, so asserting a pyspark path for it would be
false.
This is the same "declared but never applied" class of gap as SPARK-59032
(parsed
gRPC keepalive never applied to the channel) and SPARK-59037 (the 128 MiB max
message
size declared but unused): the intent existed, but the wiring that makes it
take
effect was missing.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]