A wrapper of the Apache Spark Connect client with additional functionalities that allow applications to communicate with a remote Managed Spark Session using the Spark Connect protocol without requiring additional steps.
pip install google-cloud-managed-spark-connectpip uninstall google-cloud-managed-spark-connectThis client requires permissions to manage Managed Spark Sessions and Session Templates.
If you are running the client outside of Google Cloud, you need to provide
authentication credentials. Set the GOOGLE_APPLICATION_CREDENTIALS environment
variable to point to
your Application Credentials
file.
You can specify the project and region either via environment variables or directly in your code using the builder API:
- Environment variables:
GOOGLE_CLOUD_PROJECTandGOOGLE_CLOUD_REGION - Builder API:
.projectId()and.location()methods (recommended)
-
Install the latest version of Managed Spark Connect:
pip install -U google-cloud-managed-spark-connect
-
Add the required imports into your PySpark application or notebook and start a Spark session using the fluent API:
from google.cloud.managed_spark_connect import ManagedSparkSession spark = ManagedSparkSession.builder.getOrCreate()
-
You can configure Spark properties using the
.config()method:from google.cloud.managed_spark_connect import ManagedSparkSession spark = ManagedSparkSession.builder.config('spark.executor.memory', '4g').config('spark.executor.cores', '2').getOrCreate()
-
For advanced configuration, you can use the
Sessionclass to customize settings like subnetwork or other environment configurations:from google.cloud.managed_spark_connect import ManagedSparkSession from google.cloud.dataproc_v1 import Session session_config = Session() session_config.environment_config.execution_config.subnetwork_uri = '<subnet>' session_config.runtime_config.version = '3.0' spark = ManagedSparkSession.builder.projectId('my-project').location('us-central1').dataprocSessionConfig(session_config).getOrCreate()
The ManagedSparkSession.builder provides a fluent API to configure the session. Below is a list of available methods:
| Method | Description |
|---|---|
config(key, value) |
Sets a Spark configuration property. |
dataprocSessionConfig(dataproc_config) |
Sets the Dataproc Session configuration object. |
dataprocSessionId(session_id) |
Sets a custom session ID for creating or reusing sessions. |
idleTtl(duration) |
Sets the idle time-to-live (idle TTL) for the session using a datetime.timedelta object. |
label(key, value) |
Adds a single label to the session. |
labels(labels) |
Adds multiple labels to the session. |
location(location) |
Sets the Google Cloud region. |
projectId(project_id) |
Sets the Google Cloud project ID. |
runtimeVersion(version) |
Sets the Managed Spark runtime version (e.g., "3.0"). |
serviceAccount(account) |
Sets the service account for the session. |
sessionTemplate(profile) |
Sets the Session Template to use. |
subnetwork(subnet) |
Sets the subnetwork URI for the session. |
ttl(duration) |
Sets the time-to-live (TTL) for the session using a datetime.timedelta object. |
Named sessions allow you to share a single Spark session across multiple notebooks, improving efficiency by avoiding repeated session startup times and reducing costs.
To create or connect to a named session:
-
Create a session with a custom ID in your first notebook:
from google.cloud.managed_spark_connect import ManagedSparkSession session_id = 'my-ml-pipeline-session' spark = ManagedSparkSession.builder.dataprocSessionId(session_id).getOrCreate() df = spark.createDataFrame([(1, 'data')], ['id', 'value']) df.show()
-
Reuse the same session in another notebook by specifying the same session ID:
from google.cloud.managed_spark_connect import ManagedSparkSession session_id = 'my-ml-pipeline-session' spark = ManagedSparkSession.builder.dataprocSessionId(session_id).getOrCreate() df = spark.createDataFrame([(2, 'more-data')], ['id', 'value']) df.show()
-
Session IDs must be 4-63 characters long, start with a lowercase letter, contain only lowercase letters, numbers, and hyphens, and not end with a hyphen.
-
Named sessions persist until explicitly terminated or reach their configured TTL.
-
A session with a given ID that is in a TERMINATED state cannot be reused. It must be deleted before a new session with the same ID can be created.
The package supports the sparksql-magic library for executing Spark SQL queries directly in Jupyter notebooks.
Installation: To use magic commands, install the required dependencies manually:
pip install google-cloud-managed-spark-connect
pip install IPython sparksql-magic-
Load the magic extension:
%load_ext sparksql_magic
-
Configure default settings (optional):
%config SparkSql.limit=20
-
Execute SQL queries:
%%sparksql SELECT * FROM your_table
-
Advanced usage with options:
# Cache results and create a view %%sparksql --cache --view result_view df SELECT * FROM your_table WHERE condition = true
Available options:
--cache/-c: Cache the DataFrame--eager/-e: Cache with eager loading--view VIEW/-v VIEW: Create a temporary view--limit N/-l N: Override default row display limitvariable_name: Store result in a variable
See sparksql-magic for more examples.
Note: Magic commands are optional. If you only need basic ManagedSparkSession functionality without Jupyter magic support, install only the base package:
pip install google-cloud-managed-spark-connectThe dataproc-spark-connect package has been renamed to google-cloud-managed-spark-connect. This is a breaking change with no compatibility shims — you need to update your code in the following places when you switch to the new package.
# Before
pip install dataproc-spark-connect
# After
pip install google-cloud-managed-spark-connectgoogle.cloud.dataproc_spark_connect is now google.cloud.managed_spark_connect, and DataprocSparkSession is now ManagedSparkSession:
# Before
from google.cloud.dataproc_spark_connect import DataprocSparkSession
spark = DataprocSparkSession.builder.getOrCreate()
# After
from google.cloud.managed_spark_connect import ManagedSparkSession
spark = ManagedSparkSession.builder.getOrCreate()If you use the Jupyter magic commands, google.cloud.dataproc_magics is now google.cloud.managed_spark_magics and DataprocMagics is now ManagedSparkMagics (the %dpip magic itself is unchanged).
If you set any of the library's own environment variables (as opposed to standard GCP ones like GOOGLE_CLOUD_PROJECT), rename the DATAPROC_SPARK_CONNECT_ prefix to MANAGED_SPARK_CONNECT_:
| Before | After |
|---|---|
DATAPROC_SPARK_CONNECT_SERVICE_ACCOUNT |
MANAGED_SPARK_CONNECT_SERVICE_ACCOUNT |
DATAPROC_SPARK_CONNECT_SUBNET |
MANAGED_SPARK_CONNECT_SUBNET |
DATAPROC_SPARK_CONNECT_AUTH_TYPE |
MANAGED_SPARK_CONNECT_AUTH_TYPE |
DATAPROC_SPARK_CONNECT_TTL_SECONDS |
MANAGED_SPARK_CONNECT_TTL_SECONDS |
DATAPROC_SPARK_CONNECT_IDLE_TTL_SECONDS |
MANAGED_SPARK_CONNECT_IDLE_TTL_SECONDS |
DATAPROC_SPARK_CONNECT_SESSION_TERMINATE_AT_EXIT |
MANAGED_SPARK_CONNECT_SESSION_TERMINATE_AT_EXIT |
DATAPROC_SPARK_CONNECT_DEFAULT_DATASOURCE |
MANAGED_SPARK_CONNECT_DEFAULT_DATASOURCE |
DATAPROC_SPARK_CONNECT_ACTIVE_SESSION_FILE_PATH |
MANAGED_SPARK_CONNECT_ACTIVE_SESSION_FILE_PATH |
Note that GOOGLE_CLOUD_DATAPROC_API_ENDPOINT and other variables naming the actual Dataproc API (not this library's own config) are unchanged.
For development instructions see guide.
We'd love to accept your patches and contributions to this project. There are just a few small guidelines you need to follow.
Contributions to this project must be accompanied by a Contributor License Agreement. You (or your employer) retain the copyright to your contribution; this simply gives us permission to use and redistribute your contributions as part of the project. Head over to https://cla.developers.google.com to see your current agreements on file or to sign a new one.
You generally only need to submit a CLA once, so if you've already submitted one (even if it was for a different project), you probably don't need to do it again.
All submissions, including submissions by project members, require review. We use GitHub pull requests for this purpose. Consult GitHub Help for more information on using pull requests.