This file is read by automated agents (security scanners, code analyzers, AI assistants) operating on this repository. It points them at the human-authored references they should consult before producing output.
Security model: SECURITY.md
Agents that scan this repository should consult SECURITY.md
for the project's threat model, in-scope / out-of-scope
declarations, and known non-findings before reporting issues.
Maintainer- and agent-facing internals live under
dev/design/ (not published on the user docs site):
dev/design/architecture.md- system design, parity strategy, no-stub ruledev/design/ci.md- the CI pipeline (golden-proto tests, parity gate)dev/design/official-tests.md- the official Apache Spark Connect test suite approachdev/design/test-harness.md- the official-test harnessdev/design/acceptance.md- acceptance criteria / definition of done
The API parity ledger is generated by scripts/gen_parity_ledger.py into
dev/parity/inventory.csv. User-facing documentation lives in docs/ and is
published to GitHub Pages.
This section is the condensed, hard-won process for bumping the tracked Spark version and re-achieving full parity fast. Do it in this order. The goal is 100% user-facing API parity (proven, not claimed) with >90% coverage on both the Rust and Python sides, no stubs, green CI.
- Prove parity, never claim it. The only acceptable evidence is the introspection diff in step 2 showing 0 user-facing gaps. "I added everything" is not evidence.
- Verify every claimed gap against real installed pyspark before implementing. Audit/analysis agents hallucinate methods that don't exist (fake Catalog methods, "missing" methods that are already present under a different name).
_pysparkis a compiled extension, not a package. You cannotfrom pyspark._pyspark.functions import X. Use attribute access:pyspark._pyspark.functions.pyfunc_make_udf.- Always full-test-compile before pushing (see step 7). A removed/renamed
core method has broken
e2e_dataframemore than once because a targeted test subset missed it. cargo llvm-covwithout a running server reports ~58% (e2e skips). With the server it is ~88%. Always setSPARK_REMOTEwhen measuring (step 6).- Parallel agents that edit the same crate / git index collide and have autonomously committed+pushed against instructions. If you delegate, give each agent a disjoint file set, tell them no git, no gh, no cargo-collision, and commit centrally yourself. Only one process may run cargo at a time.
- Never let an agent run
gh— it flips the active account (gh auth switch) and breaks your pushes. Rungh auth switch -u HyukjinKwonbefore pushing. - Commit messages end with
Co-authored-by: Isaac <no-reply@databricks.com>. PR bodies end withThis pull request and its description were written by Isaac.Never bypass git hooks (--no-verify). - Always merge PRs with the dev merge script
dev/merge_connect_rust_pr.py(nevergh pr mergeor a manualgit push). It produces the Apache-canonical squashed commit —Closes #<n>,Authored-by:/Co-authored-by:,Signed-off-by:— resolves the JIRA, and pushes to theapacheremote'smaster. It is interactive; drive it by piping the prompt answers on stdin (e.g. the PR number, theny/n/blank for the confirmations, author default, and "Push to apache?"). - JIRA version convention: always
connect-rust-x.y.z, never the core Sparkx.y.zline. When you file a JIRA for spark-connect-rust, set both the Affects Version/s and the Fix Version/s to aconnect-rust-x.y.zversion (e.g.connect-rust-0.1.0,connect-rust-4.2.0) — never a core Spark version. This keeps connect-rust work in its own release notes instead of leaking into core Spark's.dev/create_spark_jira.pyanddev/merge_connect_rust_pr.pyalready follow this; keep it consistent when editing tickets by hand. - Release umbrella JIRA: one Epic per
connect-rust-x.y.zrelease. Each release is tracked by an umbrella Epic titledspark-connect-rust connect-rust-x.y.z release(e.g.SPARK-59035forconnect-rust-4.2.0,SPARK-59102forconnect-rust-4.3.0), with Affects Version/s = the prior release and Fix Version/s = this release. Attach every JIRA that ships in the release to it with anis part of(Incorporates) issue link — the umbrellaincorporatesthe child, the childis part ofthe umbrella — not the Epic Link field. Use only that one mechanism so the two never drift apart.
- All cargo commands: prefix
CARGO_NET_OFFLINE=true CARGO_REGISTRIES_CRATES_IO_PROTOCOL=git. - Bump the version in the workspace
Cargo.toml(version = "X.Y.Z") and the per-cratespark-connect-proto/-corepath deps. Regenerate protos against the newspark/connect/*.proto. - Reference pyspark: install the matching official
pyspark/pyspark-clientinto a conda env (here:/opt/miniconda3/envs/python3.11/bin/python). This is the source of truth for the diff and for signatures. - Local Connect server for e2e + coverage: unpack
spark-X.Y.Z-bin-hadoop3-connect,./sbin/start-connect-server.shwith no--packages(the jar is bundled; offline maven fails), JAVA_HOME 17. It listens on IPv6[::]:15002, so an IPv4/dev/tcp/localhostprobe is misleading — check withlsof -iTCP:15002 -sTCP:LISTEN. - Build the Python extension (the dev machine is arm64 host / x86_64 conda
python, hence the explicit target — adjust per machine; CI is Linux and emits
target/release/lib_pyspark.so):Rebuild + re-copy after every Rust change you want to test from Python.cargo build -p pyspark-rs --release --target x86_64-apple-darwin cp target/x86_64-apple-darwin/release/lib_pyspark.dylib python/pyspark/_pyspark.so PYTHONPATH=$(pwd)/python /opt/miniconda3/envs/python3.11/bin/python # run the drop-in
Two Python processes, then a set-diff per class. Reference = official conda pyspark connect classes; ours = the built extension.
# official: dir() of each connect class
/opt/miniconda3/envs/python3.11/bin/python dump_official.py > /tmp/official_api.json
# ours: dir() of pyspark._pyspark.<Class>
PYTHONPATH=$(pwd)/python /opt/miniconda3/envs/python3.11/bin/python dump_ours.py > /tmp/ours_api.json
# diff official-minus-ours per class → the gap list
Cover these classes: DataFrame, Column, SparkSession, SparkSession.Builder,
DataFrameReader, DataFrameWriter, DataFrameWriterV2, GroupedData,
Catalog, Window, WindowSpec, Observation (add any new connect classes the
new version introduces).
Justified non-gaps (absent on purpose because we use a native Rust transport,
not the Python gRPC client): Column.to_plan(session) and SparkSession.client.
Anything else in the diff is a real gap to close. Re-run the diff until it shows
only these.
For every new/changed API, check the signature AND behavior against the
version's source, not just the name:
https://github.com/apache/spark/tree/vX.Y.Z/python/pyspark (and the connect
subtree python/pyspark/sql/connect/). Implement on both sides:
- Rust core (
crates/spark-connect/src/*.rs) — the real behavior / plan or proto construction. Add the method if the core lacks it. - PyO3 wrapper (
crates/pyspark-rs/src/*.rs) —#[pymethods]delegating to core, with#[pyo3(name = "camelCase", signature = (...))]. - Python shim (
python/pyspark/...) only when the class is not a direct_pysparkre-export.
Gotchas:
- Reader/Writer format methods need the FULL explicit pyspark kwargs, not a
**optionscatch-all (csv has ~35 named params). Each named kwarg → an option. py.evalevaluates expressions only — a class/defstatement raises SyntaxError. Put helper classes in a shim module and import them, or use the registration-object pattern (spark.udf().register,spark.udtf).SparkSession.Builder()needs a#[new]on the builder pyclass, plus a#[classattr] Builderexposing the class.- Properties (isStreaming, storageLevel, sparkSession, session_id, is_stopped)
are
#[getter], not methods. - State that must be observable (e.g.
is_stopped) needs real backing in core (anArc<AtomicBool>shared acrossclone()), not a hardcoded return.
- Every
DataTypemust be picklable:__reduce__→(pyspark.sql.types._parse_datatype_json_string, (json,)). That function is session-free and present in both our shim and the official worker's pyspark, so return types round-trip to official UDF workers. _parse_datatype_json_stringmust be fully recursive: handle the nested object forms (struct/array/map) as well as the string forms (atomics,decimal(p,s),char(n),varchar(n),time(n),interval …,calendarinterval). A non-recursive version silently breaks nested pickling.- Struct-field metadata values must be stored as raw strings on parse
(not
serde_json::Value::to_string, which re-quotesv→"\"v\""). Keep the JSON and proto paths consistent — seetypes_roundtrip.rs. - New DataTypes: add proto
to_proto/from_proto,json/from_json,simple_string/type_name,#[new]+ picklability, and a variant totests/types_roundtrip.rs(proto round-trip, json idempotence, accessors). - Known lossy case to watch: Geometry/Geography JSON encodes a fixed CRS token, not the numeric srid, so JSON does not round-trip (proto does).
Some python/pyspark/ dirs are verbatim upstream copies, not fork code. Keep
them as git subtrees from apache/spark at the version tag (never master),
so they are normal tracked files (importable + packaged), carry provenance in
history, and update via git subtree pull.
Rigorously classify before touching anything (diff our dir vs the version tag):
- Rust-backed shims — KEEP OURS (anything importing
pyspark._pyspark, or an adapter):sql/*top-level,ml/connect/*,sql/streaming/{query,readwriter}.py,resource/profile.py,errors/*,pyspark/__init__.py,util.py(Connect-only trim),storagelevel.py,serializers.py(fork-adapted — verify). - Pure upstream — VENDOR via subtree:
cloudpickle,pandas(see caveat),_globals.py,loose_version.py, and audit for more.
Subtree conversion (serial, needs a clean working tree, one dir at a time):
git rm -r python/pyspark/<dir> && git commit ...
( cd <clone-of-apache/spark-at-vX.Y.Z> && git subtree split --prefix=python/pyspark/<dir> -b split-<dir> )
git subtree add --prefix=python/pyspark/<dir> <clone> split-<dir> --squash
git subtree IS available in this env (the --help man page is missing, but the
command works). git clone --depth 1 --branch vX.Y.Z --filter=blob:none --sparse
keeps the clone small; subtree split on a depth-1 clone yields one commit
(fine for --squash).
pandas caveat (measured): full upstream pyspark.pandas is pure-upstream and
safe to vendor, but it does not import against the minimal Connect drop-in —
it needs deep pyspark.sql.utils/types internals (_drop_metadata,
_parse_schema, connect.functions.builtin). Making it work is a ~500–1000 LOC
bridging project. Either keep the curated subset or fund that separately; do not
ship the full copy (it won't import).
Drift guard: scripts/check_vendored_upstream.sh + .github/workflows/vendor-drift.yml
diff the vendored dirs against the pinned tag and fail CI on drift. When bumping
versions, update SPARK_TAG in the script and re-run git subtree pull for each
vendored dir, then extend VENDORED_PATHS for any newly-vendored dir.
Completeness audit: compare against the official client packaging manifest
python/packaging/client/setup.py (its connect_packages/module list) to find
files the pyspark-client distribution ships that we neither reimplement in Rust
nor vendor (e.g. errors/exceptions/connect.py, version.py, testing utils).
Vendor the ones we don't have a Rust implementation for.
- Rust core (measures
apache-spark-connectonly, not pyspark-rs):WithoutSPARK_REMOTE=sc://localhost:15002 cargo llvm-cov --no-cfg-coverage \ -p apache-spark-connect --summary-onlySPARK_REMOTEthe e2e tests skip and you see ~58% — always set it. Biggest pools to target when short:dataframe.rs,plan.rs,session.rs,streaming.rs,types.rs. Prefer offline unit tests that build a session (the gRPC channel connects lazily, so proto/plan construction needs no server) and server-gatedtests/e2e_*.rs(they checkstd::env::var("SPARK_REMOTE")). - Python drop-in:
python/.coveragercscopes to the hand-written shim (omits vendored/generated/Rust-only). Run:Mark genuinely-unreachable defensivecd python && PYTHONPATH=$(pwd) python -m coverage run --rcfile=.coveragerc \ -m pytest tests/test_dropin_offline.py && python -m coverage report --rcfile=.coveragercexcept ImportErrorfallbacks# pragma: no cover. - Both numbers feed the badges via
.github/workflows/coverage.yml.
Run all of these green before pushing:
cargo test -p apache-spark-connect --release --lib
cargo test -p apache-spark-connect --release --no-run # compiles ALL integration tests
cargo test -p apache-spark-connect --release --features wasm-udf --no-run
cargo fmt --all --check # BLOCKING CI gate
cargo clippy --workspace --all-targets # advisory (|| true in CI), but fix errors
Then re-run the step-2 diff (0 user-facing gaps) and the parity gate
scripts/run_official_tests.py (runs the official connect suite through our Rust
transport; FLAKY_FILES = retry, not skip). Only then
gh auth switch -u HyukjinKwon and push.
Run the NEW version's own official test cases — they change every release. Do not rely on the previous version's suite. When bumping to vX.Y.Z:
- Pull vX.Y.Z's official Python connect tests from
apache/sparkat the tag —python/pyspark/sql/tests/connect/**, pluspython/pyspark/{ml,pandas}/tests/…and the sharedpython/pyspark/testing/**helpers those tests import (vendor the new/changed testing utilities per step 5). New APIs ship with new test files; behavior changes ship as edited assertions. - Run them through our transport with
scripts/run_official_tests.pyand diff the pass/fail set against the reference (official pyspark running the same files). A test that passes on reference but fails on ours is a real behavior gap to fix, not a flake. RefreshFLAKY_FILES/ any skip list for the new suite (re-justify every skip; a skip is a hidden gap). - Record the resulting pass/fail delta in
dev/parity/so the next upgrade starts from a known baseline instead of rediscovering it.
- The packer is a Python script embedded in the crate
(
crates/spark-connect/src/wasm_packer.py, pulled in withinclude_str!). The Rust client runs it aspython -c <src>(so it executes as__main__), which makes cloudpickle serialize theWasmScalarUDFrunner by value into the UDF command. Because it runs as__main__there is no separatepyspark_wasm_udfpip package: the client needs only a Python interpreter withpyspark(for its vendoredcloudpickle+pyspark.sql.types._parse_datatype_json_value), and the executors need onlywasmtime. If the packer's ABI/codec changes, keep it byte-compatible withspark_connect::wasm_udf::AbiTypeand the build-time runtime inapache-spark-connect-build. - Registration mirrors pyspark:
spark.udf().register("name", &udf)(a session-sideUdfRegistrationhandle), then callable from SQL. - Removing Python on the client entirely would mean pre-baking the cloudpickled runner template per Python version and emitting the varying bytes with a Rust pickle writer; executor Python+wasmtime is intrinsic to Spark and cannot be removed.
- Bump version + regen protos; get vX.Y.Z conda pyspark + local Connect server.
- Run the introspection diff → per-class gap list (step 2) AND the module/package audit (step 10) — the class diff alone misses whole modules.
- Implement each gap in core + PyO3 (+ shim), matching signature & behavior (step 3).
- New DataTypes → picklability + json/proto round-trip +
types_roundtrip.rs(step 4). git subtree pullvendored dirs to the vX.Y.Z tag; update drift-checkSPARK_TAG(step 5).- Re-audit the packaging manifest for newly-shipped client files (step 5 / step 10).
- Coverage >90% both sides (step 6).
- Pre-push discipline: compile-all-tests, fmt, clippy, parity gate, prove-diff (step 7).
- Push (after
gh auth switch), monitor CI, fix.
The completeness contract: every module in the pyspark-client connect_packages
list (python/packaging/client/setup.py, at the version tag being tracked) must either be
vendored verbatim (pure-Python modules) OR have equivalent Rust-backed logic exposed under
the same import path. There is no third option — a connect_packages entry that is
neither present nor Rust-equivalent is a parity gap. This includes the non-obvious ones:
pyspark.testing (a public API — assertDataFrameEqual/assertSchemaEqual and the test
base classes), pyspark.pandas, pyspark.sql.pandas, pyspark.sql.plot, pyspark.logger,
pyspark.errors, pyspark.ml.*, pyspark.pipelines, etc. The only intentional exception is
the pyspark.sql.connect.* gRPC-client internals, which the Rust transport replaces — and
even those get thin re-export shims where vendored/official code (and official tests) import
from them. Enumerate connect_packages at the tracked tag and reconcile the whole list.
Parity auditing must audit packages and top-level modules, not just connect classes
(DataFrame/Column). The class-diff alone silently misses whole packages the pyspark-client
wheel ships. Every version upgrade:
- Diff the official
pyspark-clientpackaging manifest against your module inventory. Triage every missing item: client-API (implement), classic/deprecated (N/A), or server-side (N/A). Write the triage down to avoid rediscovering gaps later. - Audit top-level
pyspark/*.pyfiles andpyspark/__init__.pyexports against the manifest. Connect-N/A files (rdd,context, etc.) are distinct from public client APIs (conf,version, task context). - When auditing, check the module attributes and dunders by hand — the class introspection diff does NOT catch missing properties, dunder methods, or submodules.
When vendoring pure-Python modules (e.g. errors, logger, client-side pipelines):
- Trace the full import closure — vendoring one module may pull dependencies you
did not anticipate (e.g.,
pyspark.loggerwas itself a missing module imported bypyspark.errors). Import-time dependencies matter; runtime captures do not. - Do NOT hand-write stubs. Stub APIs (missing attributes, broken behavior) pass local
testing but fail official test suites and mis-render errors/messages. Vendor the
real module byte-for-byte, excluding only deep dependencies you cannot satisfy
(e.g.,
py4j/grpcimports in specific submodules). - Run the official tests to confirm import-time parity — official test suites will surface missing attributes the class diff and hand-written shims both miss.
For every method matching a pyspark name:
- Signature must match: default arguments, kwarg names, types. A method using
default-argument semantics (like
conf.get(key, default=None)that returns the default instead of raising when a key is absent) must route to the correct server operation. - Members documented as properties must be Python properties (
#[getter]in PyO3), not methods. The introspection diff does NOT catch property-vs-method bugs; check by hand against official docs. - Apply ergonomics sweeps consistently across ALL analogous APIs, not just the first
instance. Grep every matching arg type when a pattern emerges (e.g., string-name
lists should all take
impl IntoIterator<Item = impl Into<String>>).
- Every
DataTypemust be picklable: the__reduce__method returns(pyspark.sql.types._parse_datatype_json_string, (json,))to round-trip through official UDF workers. The_parse_datatype_json_stringfunction must handle the full recursive case (nestedstruct/array/mapforms and all atomic forms, includingdecimal,char,varchar,time,interval). Non-recursive parsers silently break nested pickling. - Struct-field metadata must preserve raw JSON, not re-quote it during parse. Keep the JSON and proto paths consistent.
df.schemaanddf.columnsmust be properties returning the correct Python type (StructType/list, not a bare DataType or method). Usedata_type_to_pyto materialize coreDataTypeobjects into their concrete Python types with the correct MRO.StructTypeandStructFieldmust expose all documented accessors and dunders (fields, names,__getitem__,__len__, etc.). The introspection diff does NOT catch missing attributes; check by hand.PyDataTypeneeds a#[new]to allow pure-Python subclasses (UserDefinedType). UDT is inherently Python and must cloudpickle the concrete class, not be lowered to Rust.
- Vendored pure-Python dirs: use
git subtreepinned to the upstream version tag (never master) so they carry provenance, are importable, and update viagit subtree pull. Classify each dir as Rust-backed shim (keep ours) or pure upstream (vendor verbatim). - Proto can intentionally diverge: the fork may extend the proto with fork-specific
fields. On a version bump, verify new upstream fields are covered with a
>-only diff (no lines from the tag should be absent from ours). The proto is exempt from vendor drift checks because it is a superset, not a byte-identical copy.
- Cloudpickled-and-server-executed APIs stay Python — do NOT rewrite them in Rust. UDFs,
pandas_udf,UserDefinedType,DataSource, andStatefulProcessor(+ ValueState/ListState/ MapState/StatefulProcessorHandle and the state clients) are all serialized with cloudpickle and executed on the server's Python worker (real pyspark). The user subclasses them in Python; the client only imports/cloudpickles them. A Rust class would (a) break the cloudpickle round-trip (the server reconstructs by module path and expects the Python class) and (b) never run on the client anyway, so it buys nothing. What IS Rust-backed is the client-side METHOD that builds the proto carrying the pickled object:registerDataSource,GroupedData.transformWithState[InPandas],createDataFrame(..., udf), etc. General rule: proto-building → Rust; cloudpickled Python code that runs server-side → vendored Python. - Row materialization: PyO3 cannot subclass
tuple, so ourRowis not a tuple subclass. Some pandas-on-Spark code assumes a tuple-Rowwith mutable__fields__. Full fidelity likely requires adopting the upstream Python tuple-Rowat collect boundaries. - Streaming listeners and progress holders: Keep progress holders (
StreamingQueryProgress,StateOperatorProgress, etc.) as faithful pure-Python implementations parsed from server JSON, not Rust-backed. The listener bus (subscribe/dispatch) is genuinely Rust-backed; the progress data structure is not. - The pandas-parity gate runs OFFICIAL pyspark over our TRANSPORT — not our drop-in classes.
.github/workflows/pandas-parity.ymlputs the apache/spark v4.2.0 checkout onPYTHONPATHand loadsscripts/rust_transport_plugin.py, which monkeypatchesSparkConnectClientso official pyspark builds every proto (Column/lit/isin/createDataFrame) and only the gRPC transport is ours. So Rust Column/lit/__repr__changes do NOT move this gate — a failure here is in the transport seam or the test environment, not the drop-in. To reproduce a failure locally, drive it the same way (do NOTimport pyspark.pandasfrom our./python):PYTHONPATH=scripts:<spark-src>/python,RUST_PYSPARK_SO=...,SPARK_REMOTE=sc://localhost:15002,-p rust_transport_plugin, against a real server. Two environment requirements the tests silently assume (both cost real time):- Client timezone must equal the server's
spark.sql.session.timeZone. The tests compare against pure-pandas (tz-naive) results; a timestamp only round-trips through Spark when the two tzs match (createDataFrame converts with the session tz,lit(datetime)with the client's local tz). REAL pyspark fails identically when they differ, so ALIGN them, don't "fix" the client. Symptoms of a mismatch:Series.isin([datetime])/get_dummieson a datetime column silently return all-False/all-zeros. Use UTC on both (server-Duser.timezone=UTC+spark.sql.session.timeZone=UTC, clientTZ=UTC), NOT a DST-having zone: America/Los_Angeles aligns too but makes client-side pandas raiseNonExistentTimeError(and skew a timestampmean) on the resample/describe tests, whose data includes the 2021-03-14 02:00 spring-forward gap. UTC also matches how Apache runs its PYTHON connect tests (the runner's UTC); the pom.xmlAmerica/Los_Angelesis only for the JVM/Scala maven-surefire tests, not the python connect server. - Pin pandas to the supported range (
>=2.2,<3.0). v4.2.0 pandas-on-Spark warns pandas ≥ 3.0.0 is unsupported and has 3.0-only branches that misbehave (e.g.test_frame_loc_setitemno-ops a reordered multi-column.locassignment). Also add test-only deps the suite needs (scipy for the corr tests, pyyaml for pipelines) — matching Apache's own pandas test environment.
- Client timezone must equal the server's
- Improving the drop-in ITSELF (separate from the gate above): fixes belong in Rust or in the
compat shims — never by editing a vendored file (breaks the drift gate) or monkey-patching a
Rust class from Python. Two gotchas that cost real time:
Column.__repr__must render the expression (Column<'<expr>'>), not a constant. pandas- on-Spark'sspark_column_equalsdecides column identity by comparing repr strings (it doesrepr(left).replace(backtick, "") == repr(right)...); a generic"Column()"makes every column look equal and collapses frames to(n, 0).Expression::render()mirrors the pyspark-connect expression__repr__(infix binary ops,AS,CAST, star, ...).pyspark.sql.connect.columnre-exports the RustColumn, so itsisinstance(x, ConnectColumn)checks already pass.- Internal pandas functions dispatch by NAME, not by client reimplementation.
pyspark.sql.internal.InternalFunction(distributed_sequence_id,pandas_product,pandas_stddev, ...) routes throughconnect.functions.builtin._invoke_function_over_columns, which must build a plainUnresolvedFunction(name, cols)(viafunctions._invoke_function→ Rustpyfunc_invoke_function). The Connect server resolves these internal names. Do NOT special-case them client-side (e.g. fakingdistributed_sequence_idwith arow_numberwindow — it is wrong and non-distributed).
*colsmethods unpack a single list.df.select(["a","b"])must behave likedf.select("a","b")(pyspark unpacks a lone list/tuple arg). Handled centrally into_column_list, so it coversselect/sort/groupBy/... at once — don't re-solve per method.- Pipelines (SDP): The pure-Python surface vendors cleanly with one edit to imports. The
Connect-server path requires a dedicated Rust
PipelineCommandseam with high-level methods on the session, separate from drop-in parity work.