Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
12 changes: 11 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,12 @@ environment variables:

### Using Spark SQL Magic Commands (Jupyter Notebooks)

The package includes the [sparksql-magic](https://github.com/cryeo/sparksql-magic) library for executing Spark SQL queries directly in Jupyter notebooks.
The package supports the [sparksql-magic](https://github.com/cryeo/sparksql-magic) library for executing Spark SQL queries directly in Jupyter notebooks.

**Installation**: To use magic commands, install with the `magic` extra:
```bash
pip install dataproc-spark-connect[magic]
```

1. Load the magic extension:
```python
Expand Down Expand Up @@ -90,6 +95,11 @@ Available options:

See [sparksql-magic](https://github.com/cryeo/sparksql-magic) for more examples.

**Note**: Magic commands are optional. If you only need basic DataprocSparkSession functionality without Jupyter magic support, install the base package:
```bash
pip install dataproc-spark-connect
```

## Developing

For development instructions see [guide](DEVELOPING.md).
Expand Down
14 changes: 14 additions & 0 deletions google/cloud/dataproc_spark_connect/session.py
Original file line number Diff line number Diff line change
Expand Up @@ -1155,6 +1155,20 @@ def stop(self) -> None:
)

self._remove_stopped_session_from_file()

# Clean up SparkSession._instantiatedSession if it points to this session
try:
from pyspark.sql import SparkSession as PySparkSQLSession

if PySparkSQLSession._instantiatedSession is self:
PySparkSQLSession._instantiatedSession = None
logger.debug(
"Cleared SparkSession._instantiatedSession reference"
)
except (ImportError, AttributeError):
# PySpark not available or _instantiatedSession doesn't exist
pass

DataprocSparkSession._active_s8s_session_uuid = None
DataprocSparkSession._active_s8s_session_id = None
DataprocSparkSession._active_session_uses_custom_id = False
Expand Down
1 change: 1 addition & 0 deletions requirements-dev.txt
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
google-api-core>=2.19
google-cloud-dataproc>=5.18
ipython~=9.1
ipywidgets>=8.0.0
packaging>=20.0
pyink~=24.0
pyspark[connect]~=4.0.0
Expand Down
8 changes: 6 additions & 2 deletions setup.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,11 +30,15 @@
install_requires=[
"google-api-core>=2.19",
"google-cloud-dataproc>=5.18",
"IPython>=7.4.0",
"packaging>=20.0",
"pyspark[connect]~=4.0.0",
"sparksql-magic>=0.0.3",
"tqdm>=4.67",
"websockets>=14.0",
],
extras_require={
"magic": [
"IPython>=7.4.0",
"sparksql-magic>=0.0.3",
],
},
Comment thread
fangyh20 marked this conversation as resolved.
Outdated
)
16 changes: 16 additions & 0 deletions tests/integration/test_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -546,6 +546,14 @@ def test_session_id_validation_in_integration(
@pytest.mark.parametrize("auth_type", ["END_USER_CREDENTIALS"], indirect=True)
def test_sparksql_magic_library_available(connect_session):
"""Test that sparksql-magic library can be imported and loaded."""
pytest.importorskip(
Comment thread
fangyh20 marked this conversation as resolved.
"IPython", reason="IPython not available (install with magic extra)"
)
pytest.importorskip(
"sparksql_magic",
reason="sparksql-magic not available (install with magic extra)",
)

from IPython.terminal.interactiveshell import TerminalInteractiveShell

# Create real IPython shell
Expand All @@ -572,6 +580,14 @@ def test_sparksql_magic_library_available(connect_session):
@pytest.mark.parametrize("auth_type", ["END_USER_CREDENTIALS"], indirect=True)
def test_sparksql_magic_with_dataproc_session(connect_session):
"""Test that sparksql-magic works with registered DataprocSparkSession."""
pytest.importorskip(
"IPython", reason="IPython not available (install with magic extra)"
)
pytest.importorskip(
"sparksql_magic",
reason="sparksql-magic not available (install with magic extra)",
)

from IPython.terminal.interactiveshell import TerminalInteractiveShell

# Create real IPython shell (DataprocSparkSession is already registered globally)
Expand Down