Skip to content
Closed
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
59 changes: 51 additions & 8 deletions docs/spark.md
Original file line number Diff line number Diff line change
Expand Up @@ -76,8 +76,31 @@ cursor = conn.cursor(session_idle_timeout_minutes=60,
})
```

The session is not terminated until the close method of the cursor is called.
You can use the context manager to automatically call the close method.
#### Session ownership

A session started by a cursor (no `session_id`) is *client-owned*. When several
cursors of the same connection share such a session (by passing its id to later
cursors), the connection reference counts the owners and the session is
terminated only when the last owner closes.

A session passed via `session_id` that the connection did not start is
*borrowed* (for example an existing session handed to PyAthena by a data
platform). Closing a borrowing cursor, or exiting its context, never
terminates that session; the caller keeps managing it.

Regardless of ownership, closing a cursor stops the calculations it started
if they have not reached a terminal state, and waits for their terminal
state. Stop and terminate requests are idempotent: concurrent execute,
cancel, timeout and close paths converge safely, and a failed termination
raises an error instead of pretending the close succeeded; calling
`close()` again reattempts the unfinished cleanup. Poll results and
stdout/stderr reads that arrive after close has started are not published.

Closing the connection (`conn.close()`, or `await conn.aclose()` for the
native asyncio connection) applies the same rules to its open Spark cursors.

The cursor's client-owned session is terminated when the last owner closes
it. You can use the context manager to automatically call the close method.

```python
from pyathena import connect
Expand All @@ -90,6 +113,18 @@ with conn.cursor() as cursor:
...
```

Borrowing a session keeps it running after the cursor closes:

```python
from pyathena import connect
from pyathena.spark.cursor import SparkCursor

# `session_id` is managed outside PyAthena (e.g. by a data platform).
with conn.cursor(SparkCursor, session_id=session_id) as cursor:
cursor.execute("...")
# The session is still alive here; only this cursor's calculations were stopped.
```

### Spark DataFrames

The Spark DataFrames code in the sample notebook that can be enabled
Expand Down Expand Up @@ -326,13 +361,16 @@ from pyathena.spark.async_cursor import AsyncSparkCursor
conn = connect(work_group="YOUR_SPARK_WORKGROUP", cursor_class=AsyncSparkCursor)
with conn.cursor() as cursor:
calculation_id, future = cursor.execute("""spark.sql("SELECT * FROM many_rows")""")
cursor.cancel(calculation_id) # The cancel method future object returns nothing.
# It is better not to get the result of cursor execution.
# Because it will be blocked until the session is terminated.
# future.result()
cursor.cancel(calculation_id).result()
# The calculation transitions to CANCELED; the session stays usable.
calculation_execution = future.result()
```

NOTE: Currently it appears that the calculation is not canceled unless the session is terminated.
Closing the cursor stops the calculations it started and releases its
session under the same ownership rules as `SparkCursor`: a borrowed session
is left running and the last owner of a client-created session terminates
it. In-flight futures resolve with an exception instead of returning data
that arrived after close started.

(aio-spark-cursor)=

Expand All @@ -357,7 +395,12 @@ async with await aio_connect(work_group="YOUR_SPARK_WORKGROUP",
print(await cursor.get_std_out())
```

The cursor supports the async context manager for automatic session termination:
The cursor supports the async context manager. Session and calculation
ownership follows the same rules as `SparkCursor`: closing a borrowing
cursor stops its calculations but leaves the caller-managed session
running, while the last owner of a client-created session terminates it.
Closing the connection with `await conn.aclose()` (or the async context
manager) applies the same rules:

```python
import asyncio
Expand Down
60 changes: 52 additions & 8 deletions pyathena/aio/connection.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,12 @@
from __future__ import annotations

import asyncio
import inspect
from typing import Any

from pyathena.aio.cursor import AioCursor
from pyathena.connection import Connection
from pyathena.error import OperationalError, ProgrammingError


class AioConnection(Connection[AioCursor]):
Expand All @@ -14,13 +16,14 @@ class AioConnection(Connection[AioCursor]):
and provides ``create()`` for non-blocking initialization.

Example:
>>> async with await AioConnection.create(
... s3_staging_dir="s3://bucket/path/",
... region_name="us-east-1",
... ) as conn:
... async with conn.cursor() as cursor:
... await cursor.execute("SELECT 1")
... print(await cursor.fetchone())
>>> async def example():
... async with await AioConnection.create(
... s3_staging_dir="s3://bucket/path/",
... region_name="us-east-1",
... ) as conn:
... async with conn.cursor() as cursor:
... await cursor.execute("SELECT 1")
... print(await cursor.fetchone())
"""

def __init__(self, **kwargs: Any) -> None:
Expand All @@ -46,8 +49,49 @@ async def create(
"""
return await asyncio.to_thread(cls, **kwargs)

async def aclose(self) -> None:
"""Close the connection and its open Spark cursors asynchronously.

Applies the same ownership rules as the synchronous
:meth:`Connection.close`: each cursor stops its unfinished
calculations, borrowed sessions survive, and the last owner of a
client-created session terminates it. Failures are surfaced and
``aclose()`` stays retryable.
"""
errors: list[Exception] = []
for cursor in self._spark_open_cursors():
close = cursor.close
try:
if inspect.iscoroutinefunction(close):
await close()
else:
await asyncio.to_thread(close)
except Exception as e:
errors.append(e)
if errors:
raise OperationalError(
"Failed to close one or more Spark cursors. Remote calculations or "
"sessions may still be running; retrying connection.aclose() reattempts "
"the unfinished cleanup."
) from errors[0]

def close(self) -> None:
"""Synchronous close is only supported without asyncio Spark cursors.

Native asyncio cursors cannot be driven from synchronous code inside
a running event loop. Use :meth:`aclose` (or the async context
manager) when :class:`~pyathena.aio.spark.cursor.AioSparkCursor`
cursors are open.
"""
if any(inspect.iscoroutinefunction(cursor.close) for cursor in self._spark_open_cursors()):
raise ProgrammingError(
"This connection has native asyncio cursors; "
"use `await connection.aclose()` to close it."
)
super().close()

async def __aenter__(self) -> AioConnection:
return self

async def __aexit__(self, exc_type, exc_val, exc_tb) -> None:
self.close()
await self.aclose()
Loading