Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .github/workflows/integration-tests.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,9 @@ jobs:
run: |
pip install -r requirements-dev.txt
pip install -r requirements-test.txt
# --no-deps keeps pip from pulling the full pyspark distribution, which
# would shadow the lightweight pyspark-client package.
pip install --no-deps sparksql-magic>=0.0.3

- name: Authenticate to Google Cloud
uses: google-github-actions/auth@7c6bc770dae815cd3e89ee6cdf493a5fab2cc093 # v3
Expand Down
35 changes: 34 additions & 1 deletion .github/workflows/tests.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -49,4 +49,37 @@ jobs:
pip install -r requirements-test.txt

- name: Run unit tests
run: python -m pytest tests/unit/ -v --tb=short -n auto
run: python -m pytest tests/unit/ -v --tb=short -n auto

local-spark:
name: Run local Spark tests
runs-on: ubuntu-latest

steps:
- name: Checkout code
uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1

- name: Setup Python
uses: actions/setup-python@5fda3b95a4ea91299a34e894583c3862153e4b97 # v7.0.0
with:
python-version: "3.12"

- name: Cache pip dependencies
uses: actions/cache@55cc8345863c7cc4c66a329aec7e433d2d1c52a9 # v6.1.0
with:
path: ~/.cache/pip
key: ${{ runner.os }}-pip-local-spark-${{ hashFiles('requirements-dev.txt', 'requirements-test.txt', 'requirements-local-spark.txt') }}
restore-keys: |
${{ runner.os }}-pip-local-spark-
${{ runner.os }}-pip-

- name: Install dependencies
run: |
pip install -r requirements-dev.txt
pip install -r requirements-test.txt
pip install -r requirements-local-spark.txt

# These tests exercise the local Spark handoff only and need no GCP
# credentials, so they run here rather than in the integration suite.
- name: Run local Spark tests
run: python -m pytest tests/integration/test_session.py -v --tb=short -k test_create_local_spark_session
25 changes: 21 additions & 4 deletions DEVELOPING.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,17 @@ pip install -r requirements-dev.txt
pip install -r requirements-test.txt
```

Tests that need a local Spark runtime are skipped unless the full `pyspark`
distribution is installed on top:

```sh
pip install -r requirements-local-spark.txt
```

Install it only when you need those tests. The unit and integration suites are
meant to run against `pyspark-client` so they keep exercising the dependency
set we ship.

# Linting/formatting

We use `pyink` to lint/format the code. To apply changes to your local
Expand Down Expand Up @@ -44,18 +55,24 @@ env \
To run tests with magic functionality, install the required dependencies manually:

```sh
pip install .
pip install IPython sparksql-magic
pip install '.[client]'
pip install IPython
pip install --no-deps sparksql-magic
```

`sparksql-magic` declares a dependency on the full `pyspark` distribution.
Installing it with `--no-deps` keeps `pyspark-client` in place; without it, pip
adds `pyspark` on top and the two shadow each other. Installing `.[full]`
instead is the other way to avoid that.

Then run tests as normal. Any magic-related tests will automatically detect and use the available dependencies.

## Testing without Magic Support

To run tests without the magic dependencies, simply install the base package:
To run tests without the magic dependencies, simply install the package:

```sh
pip install .
pip install '.[client]'
pytest
```

Expand Down
58 changes: 53 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,10 +7,44 @@ requiring additional steps.

## Install

This client needs a Spark distribution, and there are two to choose from. It
does not depend on either one directly, so that it works with whichever is
already installed. Pick with an extra:

```sh
# Talk to remote Managed Spark Sessions. Around 14 MB.
pip install 'google-cloud-spark-connect[client]'

# Also run Spark locally. Around 460 MB.
pip install 'google-cloud-spark-connect[full]'

# Neither: use the Spark distribution the environment already has.
pip install google-cloud-spark-connect
```

`[client]` installs
[`pyspark-client`](https://pypi.org/project/pyspark-client/), the Spark Connect
client on its own. `[full]` installs `pyspark[connect]`, which is the same
thing plus the Spark JVM jars — those jars are the entire size difference, and
none of them are needed to talk to a remote session.

Choose `[full]` if you also run Spark locally, or if you depend on other
packages that expect the full `pyspark` distribution. Otherwise `[client]` is
the smaller choice.

The bare install is for environments that already provide Spark, such as a
Dataproc runtime image. On its own it cannot start a session, and importing the
package tells you so. Spark 4.0 or newer is required either way.

Note that `pyspark-client` and `pyspark` both provide the `pyspark` module.
They are separate distributions, so pip will install both if asked rather than
report a conflict. Install one, not both, and switch by uninstalling the first:

```sh
pip uninstall pyspark pyspark-client
pip install 'google-cloud-spark-connect[full]'
```

## Uninstall

```sh
Expand Down Expand Up @@ -39,7 +73,7 @@ in your code using the builder API:
1. Install the latest version of Managed Spark Connect:

```sh
pip install -U google-cloud-spark-connect
pip install -U 'google-cloud-spark-connect[client]'
```

2. Add the required imports into your PySpark application or notebook and start
Expand Down Expand Up @@ -127,10 +161,24 @@ The package supports the [sparksql-magic](https://github.com/cryeo/sparksql-magi

**Installation**: To use magic commands, install the required dependencies manually:
```bash
pip install google-cloud-spark-connect
pip install 'google-cloud-spark-connect[full]'
pip install IPython sparksql-magic
```

`sparksql-magic` declares a dependency on the full `pyspark` distribution, so
it installs cleanly next to `[full]`. If you prefer `[client]`, install it
without its dependencies, otherwise pip adds `pyspark` on top of
`pyspark-client` and the two shadow each other:

```bash
pip install 'google-cloud-spark-connect[client]'
pip install IPython
pip install --no-deps sparksql-magic
```

It only imports `from pyspark.sql import SparkSession`, which `pyspark-client`
provides, so nothing is lost by skipping its dependencies.

1. Load the magic extension:
```python
%load_ext sparksql_magic
Expand Down Expand Up @@ -163,9 +211,9 @@ Available options:

See [sparksql-magic](https://github.com/cryeo/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:
**Note**: Magic commands are optional. If you only need basic ManagedSparkSession functionality without Jupyter magic support, install the package on its own:
```bash
pip install google-cloud-spark-connect
pip install 'google-cloud-spark-connect[client]'
```

## Migrating from dataproc-spark-connect
Expand All @@ -179,7 +227,7 @@ The `dataproc-spark-connect` package has been renamed to `google-cloud-spark-con
pip install dataproc-spark-connect

# After
pip install google-cloud-spark-connect
pip install 'google-cloud-spark-connect[client]'
```

### 2. Update your imports and session class
Expand Down
107 changes: 106 additions & 1 deletion google/cloud/managed_spark_connect/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,114 @@
# See the License for the specific language governing permissions and
# limitations under the License.
import importlib.metadata
import importlib.util
import warnings

from .session import ManagedSparkSession
from packaging import version

_MIN_PYSPARK_VERSION = "4.0"

_NO_SPARK_MESSAGE = (
"No Spark distribution is importable. google-cloud-spark-connect needs "
"either 'pyspark-client', for remote Managed Spark Sessions only, or "
"'pyspark', which also runs Spark locally. Install one of them with "
"'pip install google-cloud-spark-connect[client]' or "
"'pip install google-cloud-spark-connect[full]'."
)


def _installed_version(distribution):
"""Returns the installed version of a distribution, or None if absent."""
try:
return importlib.metadata.version(distribution)
except importlib.metadata.PackageNotFoundError:
return None


def _spark_import_error(exc):
"""Returns a clearer error for a missing pyspark, or None to re-raise.

Only failures to import pyspark itself are worth rewriting. Anything else
missing is a separate problem and should surface as it is.
"""
name = exc.name or ""
if name == "pyspark" or name.startswith("pyspark."):
return ImportError(_NO_SPARK_MESSAGE)
return None


def _check_pyspark_installation():
"""Checks the Spark distribution this package was installed alongside.

This package depends on no Spark distribution of its own, so that it uses
whichever one is already present. 'pyspark-client' and 'pyspark' both
provide the 'pyspark' module but are separate distributions, so pip cannot
see them as alternatives and neither can be depended on without risking a
second copy landing over the first.

That leaves three states worth reporting, since each of them otherwise
surfaces as an import error that names nothing recognizable.
"""
client_version = _installed_version("pyspark-client")
full_version = _installed_version("pyspark")

if client_version is None and full_version is None:
# Neither distribution is installed, but Spark may still be importable:
# runtime images commonly put SPARK_HOME/python on the path instead of
# installing a distribution. Only an unimportable pyspark is a problem,
# and an unmanaged one tells us no version we can go on.
if importlib.util.find_spec("pyspark") is None:
raise ImportError(_NO_SPARK_MESSAGE)
return

if (
client_version is not None
and full_version is not None
and client_version != full_version
):
warnings.warn(
f"Both 'pyspark-client' ({client_version}) and 'pyspark' "
f"({full_version}) are installed, at different versions. They "
"provide the same 'pyspark' module, so this environment holds a "
"mix of the two and imports may fail in ways that mention "
"neither. Uninstall both and reinstall only the one you need: "
"'pip uninstall pyspark pyspark-client', then "
"'pip install google-cloud-spark-connect[client]' to use remote "
"Managed Spark Sessions, or "
"'pip install google-cloud-spark-connect[full]' if you also run "
"Spark locally."
)
return

installed_version = client_version or full_version
try:
too_old = version.parse(installed_version) < version.parse(
_MIN_PYSPARK_VERSION
)
except version.InvalidVersion:
return

if too_old:
warnings.warn(
f"Spark {installed_version} is installed, but "
"google-cloud-spark-connect uses Spark Connect APIs introduced in "
f"Spark {_MIN_PYSPARK_VERSION}. Upgrade with "
"'pip install google-cloud-spark-connect[client]' or "
"'pip install google-cloud-spark-connect[full]'."
)


_check_pyspark_installation()

try:
from .session import ManagedSparkSession
except ModuleNotFoundError as e:
# The check above reads what is installed. This catches what actually
# failed to import, which covers a pyspark that is present but incomplete.
_error = _spark_import_error(e)
if _error is None:
raise
raise _error from e

old_package_names = ["google-spark-connect", "dataproc-spark-connect"]
current_package_name = "google-cloud-spark-connect"
Expand Down
6 changes: 4 additions & 2 deletions requirements-dev.txt
Original file line number Diff line number Diff line change
Expand Up @@ -4,8 +4,10 @@ ipython~=9.1
ipywidgets>=8.0.0
packaging>=20.0
pyink~=24.0
pyspark[connect]~=4.0.0
pyspark-client~=4.0.0
setuptools>=72.0
sparksql-magic>=0.0.3
# sparksql-magic declares a dependency on the full `pyspark` distribution, which
# would be installed alongside pyspark-client and shadow it. Install it without
# its dependencies instead: pip install --no-deps sparksql-magic>=0.0.3
tqdm>=4.67
websockets>=14.0
11 changes: 11 additions & 0 deletions requirements-local-spark.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
# Dependencies for the tests that need a local Spark runtime. This is what the
# [full] extra installs.
#
# Development runs against pyspark-client, which has no JVM jars and cannot
# start a local session. The Dataproc batch code path hands off to a local
# classic Spark session, so testing it needs the full distribution.
#
# Install this on top of requirements-dev.txt, never instead of it, and only
# for those tests: the unit and integration suites are meant to run against
# pyspark-client so they keep exercising the smaller of the two installs.
pyspark[connect]~=4.0.0
10 changes: 9 additions & 1 deletion setup.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,8 +32,16 @@
"google-api-core>=2.19",
"google-cloud-dataproc>=5.18",
"packaging>=20.0",
"pyspark[connect]~=4.0.0",
"tqdm>=4.67",
"websockets>=14.0",
],
# The base install deliberately names no Spark distribution, so it works
# with whichever one the environment already has. 'pyspark-client' and
# 'pyspark' both provide the 'pyspark' module but are separate
# distributions, so depending on either would install a second copy over
# the one already present. These extras are shorthand for picking one.
extras_require={
"client": ["pyspark-client~=4.0.0"],
"full": ["pyspark[connect]~=4.0.0"],
},
)
10 changes: 10 additions & 0 deletions tests/integration/test_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,9 +33,18 @@
)
from pyspark.errors.exceptions import connect as connect_exceptions
from pyspark.sql.types import StringType
from pyspark.util import is_remote_only

_SERVICE_ACCOUNT_KEY_FILE_ = "service_account_key.json"

# The library depends on pyspark-client, which has no JVM jars and therefore
# cannot start a local Spark session. Tests that need one only run when the
# full pyspark distribution is installed instead.
requires_local_spark = pytest.mark.skipif(
is_remote_only(),
reason="requires the full pyspark distribution (a local Spark runtime)",
)


@pytest.fixture(params=[None, "3.0"])
def image_version(request):
Expand Down Expand Up @@ -777,6 +786,7 @@ def local_spark_session():
session.stop()


@requires_local_spark
def test_create_local_spark_session(batch_workload_env, local_spark_session):
"""Test creating a local Spark session."""
from pyspark.sql import SparkSession as PySparkSession
Expand Down
Loading
Loading