[SPARK-59224] Add an end-to-end workflow running the gateway on kind with a real PySpark client ### What changes were proposed in this pull request? The `deploy/examples/e2e-smoke` walkthrough exercises the whole deployment path — gateway image, Helm chart, Kubernetes Endpoints discovery, session affinity, audit and metrics — but only by hand. Nothing in CI covered it, so a break in the chart, the Dockerfile or the K8s pool would only surface when someone next ran the walkthrough manually. This adds an `E2E` workflow that automates it: 1. build the gateway image 2. create a kind cluster, load the image into it 3. deploy two `apache/spark:4.0.0` Spark Connect servers 4. install the gateway with the Helm chart + the example values 5. assert the gateway discovered **both** backends (K8s Endpoints watch) 6. run the repo's own `test/integration/client_smoke.py` through a port-forward 7. assert `scg_backend_pool_size`, the `ExecutePlan` counters, and that the audit log recorded `ExecutePlan` `kind` and `helm` are installed with plain `curl` at pinned versions rather than third-party actions — both to stay clear of the ASF GitHub Actions allowlist and to keep the versions explicit. The only action used is `actions/checkout`. On failure the job dumps pod state and gateway/Spark logs; the kind cluster is always deleted. **Cost on a hosted runner — measured, not estimated:** **6m22s**, broken down as `docker build` 208s (55%), PySpark install + client 57s, `kind create cluster` 41s, Spark backends ready (~700 MiB image pull) 43s, Helm install + gateway ready 11s, tool install and assertions ~14s. It runs in parallel with the existing `build` job (2m19s), so it adds roughly four minutes to a PR's wall clock in exchange for covering the whole deployment path. Most of that is the release build, so there is headroom later if it needs trimming (caching the Rust build layers). The PySpark client is pinned to `4.0.0` to match the `apache/spark:4.0.0` backend: pyspark 4.2's `createDataFrame` reads the SQL config `spark.sql.session.localRelationSizeLimit`, which Spark 4.0 does not have, and the server answers `SQL_CONF_NOT_FOUND`. ### Why are the changes needed? This is the only coverage of the real deployment path: the Helm chart, the Kubernetes Endpoints-watch pool, and session affinity as observed by an actual Spark client. The in-process integration tests cover the proxy logic, but nothing else verifies that the chart deploys a working gateway. ### Does this PR introduce _any_ user-facing change? No. ### How was this patch tested? Ran the entire walkthrough locally before writing the workflow. The PySpark client returned correct results including a TempView query — the meaningful check that session affinity held, since a TempView lives in one driver's memory and a misrouted follow-up RPC would fail it. `scg_backend_pool_size 2`, the RPC counters and the audit records were all as expected. This PR's own run exercises the workflow on a hosted runner. ### Was this patch authored or co-authored using generative AI tooling? Yes, co-authored with Claude Code. Closes #20 from viirya/ci-04-e2e-nightly. Authored-by: Liang-Chi Hsieh <viirya@gmail.com> Signed-off-by: Liang-Chi Hsieh <liangchi.hsieh@databricks.com>
A stateless gRPC proxy that fronts a pool of Apache Spark Connect servers, providing session affinity, multi-tenant routing, auth, and observability features the open-source server intentionally leaves out.
Rust workspace built with
tonic. Production-ready surface: JWT/OIDC auth, K8s service-watch discovery, Redis-backed multi-replica HA, per-tenant routing + rate limiting, structured audit logging, Prometheus metrics, OpenTelemetry tracing, Helm chart with sensible defaults.
sc:// traffic on :15003.ExecutePlan, AnalyzePlan, Config, AddArtifacts, Interrupt, ReattachExecute, ReleaseExecute, ReleaseSession, FetchErrorDetails, CloneSession, ArtifactStatus, GetStatus) to a chosen backend.(user_id, session_id) to the same backend for the lifetime of the session — required because open-source Spark Connect keeps SparkSession state in driver-local memory.ReattachExecute / ReleaseExecute / Interrupt by operation_id via a reverse index, so reconnecting clients reach the backend that owns the operation even if the session affinity has expired.Open-source Spark Connect ships an excellent client-server protocol but deliberately leaves multi-instance coordination out of scope. For anything beyond a single Spark driver — multi-tenant platforms, HA, fleet-level observability — you need a layer in front of the servers. See OPEN_SOURCE_SPARK_CONNECT_GATEWAY_ANALYSIS.md in the plan repo for the full motivation and competitive landscape.
ExecutePlan streams.hyper has best-in-class HTTP/2 trailing-header support, which gRPC requires.Kimahriman/spark-connect-proxy, the only existing OSS Spark-Connect-native proxy.cargo build --workspace cargo test --workspace
Start a Spark Connect server on localhost:15002 (see test/integration/README.md for a Docker one-liner).
Write config.yaml:
bind_addr: ":15003" backends: - "127.0.0.1:15002"
Run the gateway:
cargo run --bin gateway -- --config config.yaml
Point a Spark Connect client at it:
from pyspark.sql import SparkSession spark = SparkSession.builder.remote("sc://localhost:15003").getOrCreate() spark.range(10).count() # → 10
The gateway exposes a Prometheus /metrics endpoint, plus /healthz and /readyz probes, on a separate admin_addr (default :9090). Set admin_addr: null to disable.
Metric set (snake_case, scg_ prefix, label cardinality bounded):
| Metric | Type | Labels | What |
|---|---|---|---|
scg_rpcs_total | counter | rpc, code | Per-RPC totals tagged by final gRPC status code |
scg_rpc_duration_seconds | histogram | rpc | Gateway-side end-to-end duration |
scg_auth_failures_total | counter | reason | Failed auth (missing_token, invalid_token, expired, unknown_kid, unknown) |
scg_backend_pool_size | gauge | — | Current healthy-backend count |
scg_active_streams | gauge | — | In-flight streaming RPCs (ExecutePlan, ReattachExecute, AddArtifacts) |
Per-RPC structured logs include a correlation ID (rid); the same ID is forwarded to the backend via x-request-id metadata so backend logs can be joined.
Off by default. Add a tracing: block to enable OpenTelemetry span export to an OTLP/gRPC collector:
tracing: endpoint: "http://otel-collector:4317" service_name: "spark-connect-gateway" sample_ratio: 1.0 # ParentBased(TraceIdRatioBased(N)) export_timeout_secs: 10
Each RPC opens an info-level scg_rpc span carrying rpc_method, rpc_system="grpc", rpc_service, and the same scg_rid correlation ID surfaced in the JSON logs. Inbound W3C traceparent metadata becomes the parent of that span; the gateway re-injects its own traceparent (alongside x-request-id) on the outbound request, so a Spark Connect server that participates in tracing joins the same trace.
When the endpoint is omitted or the whole tracing: block is absent, the gateway runs with structured JSON logs only — no OTel SDK is initialized.
Known limitation (root spans only). Today, only RPCs where the gateway is the trace root (no inbound
traceparent) export their spans reliably end-to-end via OTLP. RPCs that arrive with a W3Ctraceparentset by an upstream caller continue to log to JSON (with the inboundscg_ridcorrelation ID), but theirscg_rpcspan is dropped before reaching the OTLP exporter due to a versioning mismatch in thetracing-opentelemetry↔opentelemetry_sdkinteraction aroundContext::with_remote_span_context. Tracking upstream — distributed-trace continuity from upstream caller → gateway → backend will return once the SDK fix lands; the gateway → backend hop already injectstraceparentcorrectly so the backend half of the trace will start working as soon as the gateway-side path is fixed.
Scrape config example for Prometheus:
scrape_configs: - job_name: spark-connect-gateway kubernetes_sd_configs: - role: pod relabel_configs: - source_labels: [__meta_kubernetes_pod_container_port_name] regex: admin action: keep
By default the gateway runs without authentication (every caller is user_id: anonymous). Production deployments configure one of three authenticators:
# Bearer-token allowlist (dev / single-team): auth: type: static tokens: - { token: "alice-secret", user_id: "alice", tenant: "team-a", groups: ["devs"] } - { token: "bob-secret", user_id: "bob" }
# JWT signed by a known IdP, verified against a local public key: auth: type: jwt algorithms: ["RS256"] issuer: "https://idp.example.com" audience: "spark-connect-gateway" key: kind: pem_file path: /etc/gateway/idp-pub.pem
# OIDC / JWKS — gateway fetches keys from the IdP, refreshes on rotation: auth: type: oidc algorithms: ["RS256"] discovery_url: "https://idp.example.com/.well-known/openid-configuration" audience: "spark-connect-gateway"
In all three cases the gateway replaces whatever user_id the client declares in UserContext with the verified identity from the authenticator — clients cannot impersonate other users.
Clients pass the credential via gRPC metadata:
# PySpark Spark Connect client picks up the token from the URI: spark = SparkSession.builder.remote( "sc://localhost:15003/;token=alice-secret" ).getOrCreate()
Authenticating clients at the gateway only helps if clients cannot dial the backends directly. For that, start the backends with spark.connect.authenticate.token (Spark 4.0+) and give the token only to the gateway:
# Presented as `authorization: Bearer …` on every gateway→backend # request; per-tenant pools can override it. Backends then reject # clients that bypass the gateway with UNAUTHENTICATED. backend_token: kind: env # inline | env | file name: SCG_BACKEND_TOKEN
See Enforcing the trust boundary for the full story (including the complementary NetworkPolicy).
The default in-memory affinity store works only for a single gateway replica — every replica keeps its own (user_id, session_id) -> backend table, so a client that lands on replica B after replica A bound it to backend X gets re-bound to a different backend, breaking the Spark Connect per-driver session invariant.
To run the gateway with replicas > 1 (e.g. behind a Kubernetes Service for HA), point the gateway at a Redis instance:
bind_addr: ":15003" backends: ["spark-connect-1:15002", "spark-connect-2:15002"] affinity_store: type: redis url: "redis://redis-cluster:6379" # supports redis://:pw@host:6379/2 key_prefix: "scg-prod" # default "scg" session_ttl_secs: 3600 # default 1h, refreshed on reads op_ttl_secs: 900 # default 15min
Both stickiness invariants are preserved across replicas:
bind_session_if_absent uses Redis SET … NX EX — when two replicas race on the same session, exactly one wins; the loser reads the winner's value back and routes to the same backend.ReattachExecute / ReleaseExecute / Interrupt) lives in {prefix}:o:{op_id}, so a client reconnecting through a different replica still reaches the original driver.If Redis becomes unreachable, the gateway logs warn! per failed operation and degrades to pool-only routing — sessions land on whatever backend the pool picks each time, which is exactly the single-replica in-memory behaviour. Service remains available; HA stickiness recovers as soon as Redis does.
affinity_store defaults to type: memory (no entry needed), matching the single-replica baseline.
crates/proxy/examples/ha_smoke.rs spins up two real SparkConnectProxy instances backed by one Redis and a shared pool of fake backends, then drives RPCs through different replicas to prove three invariants:
ReattachExecute(op_id, session_id="different") arriving at replica B (after A is gone) still reaches the backend that ran the original ExecutePlan.Run with a Redis listening on :6399 (or override via REDIS_URL):
redis-server --port 6399 --daemonize yes cargo run -p scg-proxy --example ha_smoke
Exits zero on success; panics with a descriptive assertion message on any failure, so a CI script can wrap it directly.
A Helm chart is shipped at deploy/helm/scg/. Quickstart:
helm install scg ./deploy/helm/scg \ --namespace spark-connect --create-namespace
Defaults give you 2 gateway replicas + a bundled Redis StatefulSet (AOF-persisted), with a static backend list pointing at spark-connect-{1,2}.svc.cluster.local:15002. Switch to K8s service-watch discovery, JWT auth, OTLP tracing, or an external managed Redis with a few values flips — see the chart's values.yaml and README.md for the full reference.
The chart fails template-time if you set replicaCount > 1 together with affinityStore.type: memory, since that combination silently breaks Spark Connect's per-driver session invariant — redis is the default for exactly that reason.
For day-2 operations, see:
docs/architecture.md — design overview for people reading the code: layered crate structure, request flow end-to-end, what each of the 15 crates does and why it's its own crate, key design decisions, where to look first.docs/deployment.md — from-zero deployment runbook: prerequisites, install, hardening, upgrades, uninstall.docs/multitenancy.md — multi-tenant setup guide: decision matrix, three sample configs (permissive, strict, single-tenant), migration paths.docs/observability.md — every scg_* metric explained, PromQL examples, log line anatomy, distributed tracing guide, suggested alerts.docs/runbook.md — symptom → diagnosis → fix for the failures you'll actually hit (CrashLoopBackOff, /readyz 503, broken stickiness, Redis outage, auth failure spikes, K8s RBAC, latency regressions).docs/perf-baseline.md — performance baseline numbers (unary throughput, streaming concurrency, health-check overhead, drain semantics under load) reproducible via the load harness in crates/proxy/examples/load.rs.End-to-end walkthroughs (each runs on a fresh kind cluster):
deploy/examples/e2e-smoke/ — Helm install, real Spark Connect servers, PySpark client through the gateway. Verifies session affinity and audit events. Start here.deploy/examples/e2e-scale-test/ — scale the backend Deployment up and down at runtime; verify the K8s Endpoints watcher in pool-k8s picks up the changes without a gateway restart.deploy/examples/e2e-multi-replica-redis/ — two gateway replicas backed by the bundled Redis StatefulSet; verify affinity bindings survive a gateway pod restart because they live in Redis, not in pod memory.deploy/examples/e2e-trust-boundary/ — backends enforce spark.connect.authenticate.token and only the gateway holds it; verify a direct-to-backend connection is refused with UNAUTHENTICATED while the same client succeeds through the gateway.For production-style manifests using the upstream apache/spark-kubernetes-operator, see deploy/examples/spark-connect-server/.
Once those servers (and a fronting Service) exist, point the gateway at the Service's Endpoints and let the gateway pick up backends automatically:
bind_addr: ":15003" backend_discovery: type: k8s namespace: spark-connect service_name: spark-connect port: 15002
The gateway watches the Endpoints object via kube-rs. When pods are added, removed, or restarted, the gateway's backend list updates within seconds — no kubectl rollout of the gateway, no config edit. The gateway pod needs a ServiceAccount with get, list, and watch on endpoints in the target namespace; the Helm chart at deploy/helm/scg/ wires up the ServiceAccount and Role automatically.
client (sc://) ──▶ gateway ──┬──▶ Spark Connect server #1
├──▶ Spark Connect server #2
└──▶ Spark Connect server #N
scg-store-redis so all replicas share the same affinity table.crates/ gateway/ # binary entry point proxy/ # SparkConnectService impl that forwards every RPC routing/ # SessionKey, Pool/AffinityStore traits, Router, TenantRouter store-memory/ # in-memory AffinityStore (single-replica) store-redis/ # Redis-backed AffinityStore (multi-replica HA) pool-static/ # static backend pool (round-robin) pool-k8s/ # K8s Endpoints-watch backend pool healthcheck/ # HealthAwarePool — active gRPC health probes auth/ # static-token / JWT / OIDC authenticators tenant/ # tenant resolver (from-claim / from-metadata / fixed) ratelimit/ # in-memory + Redis-backed token bucket audit/ # structured audit-event emitter observability/ # metrics, tracing, admin server config/ # YAML config loader genproto/ # tonic-generated bindings for spark.connect.* proto/spark/connect/ *.proto # vendored read-only mirror of upstream deploy/examples/ spark-connect-server/ # K8s manifests (apache/spark-kubernetes-operator) test/integration/ README.md, client_smoke.py # real PySpark E2E
The crates/genproto/build.rs script invokes tonic-prost-build on every cargo build. To force a regeneration:
cargo clean -p scg-genproto cargo build -p scg-genproto
protoc must be on $PATH (e.g. brew install protobuf).
Proxy core. Streaming forward of every SparkConnectService RPC. Session affinity pinned by (tenant, user_id, session_id) so a client's RPCs always reach the same backend driver. Operation-id reverse index so ReattachExecute / ReleaseExecute / Interrupt reach the right driver even after the session binding has expired.
Backend discovery. Static list, or K8s Endpoints watch via kube-rs — pool updates within seconds of pod changes, no gateway restart.
Multi-replica HA. Redis-backed shared affinity store. Two-step graceful shutdown drains in-flight streams before the gRPC server stops.
Auth. Pluggable: anonymous (default — trusted networks only), static token, JWT (local public key), OIDC (remote JWKS / discovery). Verified Identity overwrites the client-supplied UserContext.user_id on forward.
Multi-tenancy. Tenant resolver reads from the auth claim, a gRPC metadata header, or a fixed string. Per-tenant backend pool overrides with UseDefault / Reject policies for unknown tenants. Per-tenant token-bucket rate limiting with optional per-user sub-bucket; in-memory or Redis-shared.
Audit + observability. Structured audit events (session.create, session.release, auth.failure, rpc.error) with target=scg::audit. Prometheus metrics covering RPC throughput / duration / auth failures / pool size / active streams / rate-limit rejections. OpenTelemetry tracing with W3C traceparent propagation. Active gRPC health probing.
Deployment. Helm chart with sensible defaults (2-replica HA + bundled Redis); separate values overlays for K8s discovery, JWT/OIDC auth, external Redis, tracing.
Known gaps. Weighted backend selection per tenant, cold-start provisioning of new tenant pools, per-tenant warm pools, and multi-cluster discovery are roadmap items — see ROADMAP.md for the full list with current status, planned shapes, and deliberate non-goals. Distributed-trace continuity for inbound traceparent is limited by upstream tracing-opentelemetry ↔ opentelemetry_sdk plumbing; root-span traces work end-to-end.
Licensed under the Apache License, Version 2.0. The vendored Spark Connect protos under proto/spark/connect/ are themselves under the Apache 2.0 license held by the Apache Software Foundation; see NOTICE for details.
Contributions are accepted under the same license — by submitting a pull request you agree that your contribution is licensed under Apache 2.0.