Distinguish borrowed and owned Spark sessions on cursor/connection close - #787
MonsterRHD wants to merge 4 commits into
Conversation
Spark cursors previously always called terminate_session on close, even for a session explicitly provided via session_id, so exiting one cursor killed sessions still in use by other tasks. Self-created sessions and their calculations could also leak on exception/cancel paths. - Track session ownership per connection: externally provided sessions are borrowed and never terminated; sessions started by cursors are reference counted and terminated only by their last owner. - Track every calculation a cursor starts; close stops calculations that have not reached a terminal state, idempotently, and awaits their terminal state. Stop/terminate treat an already terminal calculation or terminated session as success. - Gate polling and stdout/stderr reads once close begins so late results are not published; serialized close locks keep concurrent closes safe and failed cleanup raises a retryable OperationalError instead of pretending success. - Apply the same rules from connection close: Connection.close and the new AioConnection.aclose (used by the async SQLAlchemy adapter) close tracked cursors with the same ownership semantics. - Cover sync, thread-based and native asyncio cursors with offline ownership/race/idempotency tests and update cancel integration tests.
| :class:`~pyathena.error.OperationalError` describing the unfinished | ||
| remote cleanup. | ||
| """ | ||
| with self._close_lock: |
There was a problem hiding this comment.
Self-review round 1 — behavior/failure paths
Verified concurrent close() calls serialize here: without this lock two callers could each run _release_owned_session() and decrement the connection-level owner count twice for one cursor (a count of 2 could wrongly hit the terminate branch while another owner cursor was still open). The thread lock covers the sync and thread-pool cursors; AioSparkCursor.close mirrors this with an asyncio.Lock (safe to construct in the to_thread init on Python >= 3.10, loop binds lazily). Regression: test_concurrent_close_releases_session_ownership_once (sync threads + asyncio gather).
| self._raise_if_closing() | ||
|
|
||
| def guarded() -> Any: | ||
| self._raise_if_closing() |
There was a problem hiding this comment.
Self-review round 1 — repair (resource/data boundary)
Finding: queued worker tasks only checked the closing gate after their work, so a stdout/stderr S3 read submitted before close(wait=False) could still issue the S3 download before failing (and only the result was suppressed). Repair: gate at the start of the guarded task too, so work that had not started never runs after close; the trailing gate remains for tasks whose call was already in flight. Polls already gate inside __poll. Coverage: thread-cursor race tests assert the future resolves with ProgrammingError and no data is published.
| _logger.exception("Failed to cancel calculation.") | ||
| raise OperationalError(*e.args) from e | ||
|
|
||
| def _stop_calculation(self, calculation_id: str) -> AthenaCalculationExecution: |
There was a problem hiding this comment.
Self-review round 2 — AWS operational semantics (claims audit)
Checked the idempotency claims against the current Athena API reference:
StopCalculationExecutionon a calculation already in a terminal state "succeeds but has no effect", and cancellation is best-effort (a calc can still finish COMPLETED/FAILED). The status-first guard here avoids the call entirely for known-terminal calcs; a racing terminal transition after the call is observed via the status re-check, and the wait loop accepts COMPLETED/FAILED/CANCELED rather than assuming CANCELED.TerminateSessionon an already inactive session (TERMINATED/TERMINATING/FAILED) also succeeds with no effect;InvalidRequestExceptionis therefore only treated as benign when a freshGetSessionStatusconfirms TERMINATED — otherwise the error is surfaced and the ownership claim is rearmed.
Boundedness note (consistent with the existing __poll/_wait_for_idle_session loops): the terminal-state wait here is unbounded and paced by poll_interval; it does not add extra retries beyond the existing retry_api_call wrapper. This matches the requirement to obtain a terminal state, and is the same trade-off already present in execute polling; not changed in this PR.
|
|
||
| def close(self) -> None: | ||
| self._connection.close() | ||
| await_only(self._connection.aclose()) |
There was a problem hiding this comment.
Self-review round 2 — caller compatibility
Connection.close()keeps its signature; for non-Spark connections the registry is empty and it remains a non-raising no-op, so every existing sync/SQLAlchemy caller is unaffected. With open Spark cursors it may now raiseOperationalErroron failed cleanup (the intended retryable signal).AioConnection.__aexit__now drivesaclose(); the aio test fixtures that still call the synchronousconn.close()only hold regular aio cursors, which leave the Spark registry empty and therefore do not trip the newProgrammingErrorguard. AioSparkCursor is not exposed through any SQLAlchemy dialect, but the async adapter's pool close path was smoke-tested offline throughgreenlet_spawn/await_only; fulljust test sqla-asyncremains CI-gated (needs AWS).- Public shapes preserved:
SparkCursor/AsyncSparkCursor(close(wait=False))/AioSparkCursor.close()signatures, explicitsession_idexistence + idle validation, engine configuration, idle timeout, and completed-result/stdout reads before close.
CLEAN within round-two scope; runtime AWS behavior of the cancel/borrowed-session integration tests is deferred to CI as stated in the PR body.
Address independent review findings: - Only release session ownership after this cursor's calculations were stopped, so a failed stop can never decrement another owner's count or terminate a shared session; failed close retries keep ownership intact. - Mark sessions terminating while the last owner's termination is in flight so a concurrent cursor cannot classify them as borrowed; joining a registered session now verifies liveness and rolls back its count if construction fails; the S3 client is created before ownership is taken. - Reconcile calculations whose start API returns after close converged: they are stopped best-effort and never published (sync, thread-pool and asyncio execute paths), including a synchronous gate after S3 reads and a lock-protected terminal publication. - Rearm ownership on BaseException (incl. asyncio cancellation and KeyboardInterrupt) during termination. - Fix the borrowed-session integration test to use SparkCursor.
|
Thank you for the PR. This PR changes session ownership, calculation cleanup, and connection shutdown across multiple cursor Please explain who creates and shares the Spark session in your application, who is responsible for terminating |
WHAT
Spark cursors now distinguish sessions they own from sessions they borrow, and converge calculations and sessions by ownership, across the synchronous (
SparkCursor), thread-based (AsyncSparkCursor) and native asyncio (AioSparkCursor) cursors.session_idthat this connection did not start is borrowed and closing the cursor (or exiting its context) never terminates it. A session started by a cursor is registered on the connection and reference counted; cursors on the same connection that reuse its id join the ownership, and only the last owner terminates the session.InvalidRequestExceptionracing with an idle-timeout or parallel termination). Execute, cancel, KeyboardInterrupt/CancelledError, poll failure and context exit all converge: once close begins, poll loops bail and stdout/stderr/result reads raise instead of publishing late data. Concurrent closes are serialized (thread lock andasyncio.Lock) and release the ownership count exactly once.close()raisesOperationalError, keeps the cursor registered and the ownership claim armed, so a laterclose()retries instead of silently succeeding.Connection.close()closes tracked Spark cursors with the ownership rules above.AioConnectiongainsaclose()(awaited from__aexit__and the async SQLAlchemy adapter viaawait_only); its synchronousclose()raises a clearProgrammingErrorwhile native asyncio Spark cursors are open.docs/spark.md) describe the ownership model; existingsession_idreuse, engine configuration, idle timeout, normal result/log reads, and the thread cursorclose(wait=...)API are preserved.WHY
When a data platform hands an existing Athena Spark session to PyAthena, any cursor exiting its context called
terminate_session, killing tasks on that session that were still running. Conversely, sessions created by cursors could leak along with remote calculations when execute hit an exception, cancellation or poll-failure path, and teardown races could publish results from a poll after close or fake success when termination failed.Validation
uvx ruff@0.14.14 check .,ruff format --check .,uv run mypy .: pass locally.markdownlint-cli2@0.18.1 docs/spark.md: pass locally.tests/pyathena/spark/test_session_ownership.py(26 tests, no AWS): borrowed vs owned close, last-owner termination, shared sessions, consecutive calculations, idempotent stop/terminate, terminate failure + retry, poll-failure cleanup, close/poll/read races for all three cursor implementations (incl. deterministic thread and asyncio races), concurrent close, and connectionclose()/aclose()rules.test_cancelnow verifies cancellation without terminating the session plus a follow-up calculation; added a borrowed-session integration test). These require AWS and will run in CI.greenlet_spawn/await_onlyoffline; fulljust test sqla/sqla-asyncrequire AWS and are pending CI.