Skip to content
Great engineers know 3 min read

Spark Connect

A thin client that talks to a remote Spark server over gRPC, and what it changes for applications.

You will learn

  • What Spark Connect is and how it differs from a classic session
  • How a client talks to the server
  • What you gain: stability, upgrades, lighter clients
  • What does not work over Connect

Read first

Comfortable with these? Read on.

TL;DR Spark Connect splits the client from the driver. Your application builds unresolved plans and sends them over gRPC to a Spark server, which runs them and streams results back as Arrow batches. Introduced in Spark 3.4.

Classic mode

In a classic PySpark application, your Python process starts a JVM driver locally (through Py4J) and the driver runs inside your application. The client and the driver share a fate and a version: a memory leak in your app can kill the driver, and upgrading Spark means upgrading every application at once.

Connect mode

  1. Client

    Builds an unresolved plan from DataFrame calls

  2. gRPC

    Sends the plan (protobuf) to the server

  3. Server

    A long-running driver analyses, optimises and executes it

  4. Arrow

    Results stream back as Arrow batches

PySpark · Connecting
from pyspark.sql import SparkSession

spark = SparkSession.builder.remote("sc://spark-server:15002").getOrCreate()
spark.read.table("sales").groupBy("country").count().show()

The server is started with sbin/start-connect-server.sh (default port 15002). Most DataFrame and SQL code runs unchanged; the API is the same pyspark.sql DataFrame API.

Why it matters

Stability

A client crash or out-of-memory error does not take down the shared driver, and vice versa.

Upgrades

Server and clients can be upgraded independently, within compatibility limits.

Thin clients

Clients need no JVM. Connect clients exist for Python, Scala, and community clients for Go, Rust and others. Spark 4.0 ships a lightweight Python-only client package.

It is also what powers notebooks and IDEs talking to remote clusters, and serverless offerings: the client is just a small library speaking a protocol.

What changes

  • No RDD API and no SparkContext on the client: spark.sparkContext, df.rdd and RDD-based libraries do not work. Code must use DataFrames and SQL.
  • Analysis is deferred: in classic mode a bad column name fails as soon as you write the transformation; over Connect, plans are analysed on the server, typically when an action runs or the schema is requested. Errors can surface later.
  • Schema access is a round trip: df.columns and df.schema ask the server. Calling them in a tight loop is slow.
  • Python UDFs still work: the code is serialised and runs on the server's executors, so the server needs compatible Python and libraries.
Note: Spark 4.0 makes Connect much more complete and lets the same application switch between classic and Connect modes by configuration. Check the version notes of your platform before relying on a specific API.

Common mistakes

Using df.rdd or sparkContext

Not available over Connect.

Calling df.schema in loops

Each call is a server round trip.

Assuming errors appear immediately

Analysis happens on the server, often at action time.

Key takeaways

  • Spark Connect sends unresolved plans over gRPC to a remote driver.
  • Results come back as Arrow batches.
  • It isolates clients from the driver and decouples upgrades.
  • No RDD API or SparkContext on the client; analysis errors can appear later.

Check yourself

3 questions

1. What does a Spark Connect client send to the server?

Show the answer

Unresolved logical plans over gRPC. DataFrame calls become protobuf-encoded unresolved plans.

2. Which of these does not work over Spark Connect?

Show the answer

df.rdd.map(...). The RDD API and SparkContext are not available to Connect clients.

3. In which version was Spark Connect introduced?

Show the answer

3.4. Spark Connect arrived in Spark 3.4.

Go deeper

Primary sources: Spark Connect overview