Skip to content
15 changes: 15 additions & 0 deletions docs/spark.md
Original file line number Diff line number Diff line change
Expand Up @@ -285,6 +285,18 @@ requests cancellation, waits until the calculation reaches a terminal state, and
The `state` property returns that terminal state.
If the cancellation request fails, the `KeyboardInterrupt` propagates with the error as its cause.

A `KeyboardInterrupt` while `execute()` is still starting the calculation first waits for the
[StartCalculationExecution](https://docs.aws.amazon.com/athena/latest/APIReference/API_StartCalculationExecution.html)
request to finish, and then cancels the calculation it started in the same way.
The `calculation_id` property returns that calculation's ID.
A second `KeyboardInterrupt` during this wait propagates at once without cancelling the calculation.
A cancellation request sent right after a calculation starts can occasionally have no effect, so the calculation can still end in the `COMPLETED` state.

Unless `client_request_token` is passed to `execute()`, the cursor sends a generated `ClientRequestToken` with each calculation.
A retried start request then returns the calculation that an earlier attempt started instead of starting another one.
A token passed to `execute()` must be unique for each calculation:

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round two — claims, callers, operations (head c0188ced93137cd106a416bd092affaf13c4a552)

Result: FINDINGS, repaired in c0188ce and the PR body.

Claims checked:

  • botocore 1.43.102 marks ClientRequestToken as idempotencyToken for StartQueryExecution but not for StartCalculationExecution / StartSession — confirmed from the service model.
  • Token replay returns the same calculation, including with a different CodeBlock; a different token on a busy session fails with InvalidRequestException; a stop right after the start was ignored in 1 of 6 attempts — from the Spark: an interrupt during StartCalculationExecution leaves the calculation running #841 measurement.
  • 28/33 failing on base, 80 offline passed, 95 passed / 2 skipped live with no sessions left — rerun results for this head.

Findings and repairs:

  1. The PR body said a retried start after a lost response "previously failed"; that case was inferred from the busy-session measurement, not observed. Reworded to state the measured part.
  2. Docs did not say that a caller-supplied client_request_token reused with different code silently returns the earlier calculation (measured). Added to this paragraph.

Caller compatibility: _calculate's signature and return are unchanged; execute(client_request_token=...) still sends the caller's token; kill_on_interrupt=False keeps the old control flow. AWS operator: one extra thread per execute() with kill_on_interrupt; retries of the start request now reuse one token, so they cannot create a second calculation.

Athena returns the earlier calculation for a reused token, even when the code differs.

(async-spark-cursor)=

## AsyncSparkCursor
Expand Down Expand Up @@ -498,3 +510,6 @@ async with await aio_connect(work_group="YOUR_SPARK_WORKGROUP",

With `kill_on_interrupt` enabled, which is the default, cancelling the task while `execute()` waits for the calculation
requests cancellation of the calculation, waits until it reaches a terminal state, and then raises `asyncio.CancelledError`.
Cancelling the task while `execute()` is still starting the calculation first waits for the start request to finish,
and then cancels the calculation it started in the same way.
Cancelling the task again during this wait raises `asyncio.CancelledError` at once without cancelling the calculation.
80 changes: 73 additions & 7 deletions pyathena/aio/spark/cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@

import asyncio
import logging
import uuid
from typing import Any, cast

from pyathena.aio.util import async_retry_api_call
Expand Down Expand Up @@ -133,24 +134,78 @@ async def _calculate( # type: ignore[override]
description: str | None = None,
client_request_token: str | None = None,
) -> str:
"""Start a calculation execution with ``StartCalculationExecution``.

Without ``client_request_token``, a generated token is sent, so that a
retried request returns the calculation an earlier attempt started instead
of starting another one.

With ``kill_on_interrupt`` enabled, the request is shielded from task
cancellation. On cancellation, waits for the request to finish, requests
cancellation of the calculation it started, waits for a terminal state,
stores the calculation ID and execution on the cursor, and re-raises
``asyncio.CancelledError``. Another cancellation during that wait
propagates at once.

Args:
session_id: The session ID.
code_block: The code to run.
description: The calculation description.
client_request_token: The idempotency token of the request.

Returns:
The calculation execution ID.

Raises:
asyncio.CancelledError: If the task is cancelled while starting the
calculation. A failure to start, cancel, or wait for the
calculation becomes its ``__cause__``.
DatabaseError: If the request fails.
"""
request = self._build_start_calculation_execution_request(
session_id=session_id,
code_block=code_block,
description=description,
client_request_token=client_request_token,
client_request_token=client_request_token or str(uuid.uuid4()),
)
if not self._kill_on_interrupt:
return await self.__start_calculation(request)

start = asyncio.ensure_future(self.__start_calculation(request))
try:
return await asyncio.shield(start)
except asyncio.CancelledError as cancellation:
_logger.warning("Query canceled by user.")
try:
self._calculation_id = await start
await self.__cancel_and_wait(self._calculation_id)
except Exception as e:
raise cancellation from e
raise

async def __start_calculation(self, request: dict[str, Any]) -> str:
"""Send a ``StartCalculationExecution`` request.

Args:
request: The request parameters.

Returns:
The calculation execution ID.

Raises:
DatabaseError: If the request fails.
"""
try:
response = await async_retry_api_call(
self._connection.client.start_calculation_execution,
config=self._retry_config,
logger=_logger,
**request,
)
calculation_id = response.get("CalculationExecutionId")
except Exception as e:
_logger.exception("Failed to execute calculation.")
raise DatabaseError(*e.args) from e
return cast(str, calculation_id)
return cast(str, response.get("CalculationExecutionId"))

async def __poll(self, query_id: str) -> AthenaQueryExecution | AthenaCalculationExecution:
while True:
Expand Down Expand Up @@ -195,14 +250,25 @@ async def _poll( # type: ignore[override]
raise
_logger.warning("Query canceled by user.")
try:
await self._cancel(query_id)
self._calculation_execution = cast(
AthenaCalculationExecution, await self.__poll(query_id)
)
await self.__cancel_and_wait(query_id)
except Exception as e:
raise cancellation from e
raise

async def __cancel_and_wait(self, calculation_id: str) -> None:
"""Request cancellation and store the calculation's terminal state.

Args:
calculation_id: The calculation execution ID.

Raises:
OperationalError: If the cancellation or a status request fails.
"""
await self._cancel(calculation_id)
self._calculation_execution = cast(
AthenaCalculationExecution, await self.__poll(calculation_id)
)

async def _cancel(self, query_id: str) -> None: # type: ignore[override]
request: dict[str, Any] = {"CalculationExecutionId": query_id}
try:
Expand Down
140 changes: 135 additions & 5 deletions pyathena/spark/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,14 +9,17 @@

import contextlib
import logging
import threading
import time
import uuid
from abc import ABCMeta, abstractmethod
from concurrent.futures import Future, wait
from datetime import datetime
from typing import Any, cast

import botocore

from pyathena import NotSupportedError, OperationalError
from pyathena import DatabaseError, NotSupportedError, OperationalError
from pyathena.common import BaseCursor
from pyathena.model import (
AthenaCalculationExecution,
Expand All @@ -28,6 +31,11 @@

_logger = logging.getLogger(__name__)

# How often a wait for the start request wakes up to check for Ctrl-C, so that
# a KeyboardInterrupt is raised promptly where an untimed lock wait cannot be
# interrupted by signals (Windows before Python 3.14).
_INTERRUPT_CHECK_INTERVAL = 0.1

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Repair: rename _START_WAIT_INTERVAL to _INTERRUPT_CHECK_INTERVAL and rebase onto 9a1a72d (maintainer feedback that the old name was confusing). The PR was moved back to Draft for this repair.

The series moved from ceb9ad7..bdfa630 to 9a1a72d..1466e7b. git range-diff shows commits 1–6 identical and one new commit, 1466e7b. Upstream changes between the two bases touch only .github/workflows/* and docs/testing.md; per #863, PRs now run the AWS tests on the newest Python only. No Spark source or test changed upstream, and the rebase was clean.

What changed:

  • Constant renamed; the value stays 0.1. Its single use is __wait_for_start. No other reference exists in the code, tests, docs, or PR text.
  • The comment now states its purpose first: "How often a wait for the start request wakes up to check for Ctrl-C, so that a KeyboardInterrupt is raised promptly where an untimed lock wait cannot be interrupted by signals (Windows before Python 3.14)."
  • The platform claim was checked against the Python docs for threading.Lock.acquire: "Changed in version 3.2: … can now be interrupted by signals on POSIX" and "Changed in version 3.14: … on Windows". The previous comment's "on every platform" was therefore narrowed to the case the timeout actually covers within the supported 3.10+ range.

Self-review round one (behavior): CLEAN. The rename is behavior-neutral. After the rebase, just lint passes; offline Spark tests give 113 passed, 2 skipped (plus the 14 live-fixture setup errors expected with --noconftest). A real-SIGINT probe on 3.13.1 still cancels the returned ID.

Self-review round two (claims): CLEAN. The comment's version and platform claim matches the Python docs, and the commit message says the same. just docs lint reports 0 errors.

Independent follow-up (relayed): CLEAN.

  • Reviewer: Codex CLI 0.157.0, read-only, model_reasoning_effort=high, run on a detached snapshot of 1466e7b with the range-diff, the upstream delta, and the new commit, without the PR text or prior findings. Static review; it ran no tests.
  • It confirmed the six commits are unchanged, found no interaction with the upstream CI changes, and found no remaining references to the old name.
  • It traced wait() → Event.wait() → Condition.wait() and judged the Windows before 3.14 qualification correct.



class SparkBaseCursor(BaseCursor, metaclass=ABCMeta):
"""Abstract base class for Spark-enabled cursor implementations.
Expand Down Expand Up @@ -340,14 +348,136 @@ def _poll(self, query_id: str) -> AthenaQueryExecution | AthenaCalculationExecut
raise
_logger.warning("Query canceled by user.")
try:
self._cancel(query_id)
self._calculation_execution = cast(
AthenaCalculationExecution, self.__poll(query_id)
)
self.__cancel_and_wait(query_id)
except Exception as e:
raise interrupt from e
raise

def __cancel_and_wait(self, calculation_id: str) -> None:
"""Request cancellation and store the calculation's terminal state.

Args:
calculation_id: The calculation execution ID.

Raises:
OperationalError: If the cancellation or a status request fails.
"""
self._cancel(calculation_id)
self._calculation_execution = cast(AthenaCalculationExecution, self.__poll(calculation_id))

def _calculate(

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Self-review round one — implementation behavior (base 0ba0875c72c702e0b893843d3c50aa98b4a61b8f, head 73a738e4fd01b00e3a1f947959c730867dcf0767; the later head c0188ced93137cd106a416bd092affaf13c4a552 only adds a docs sentence)

Covered: SparkBaseCursor._calculate / __start_calculation / __wait_for_start / __cancel_and_wait, AioSparkCursor._calculate / __start_calculation / __cancel_and_wait, the unchanged _poll recovery, callers SparkCursor.execute, AsyncSparkCursor.execute (starts on the caller's thread), AioSparkCursor.execute, and the new offline tests.

Result: CLEAN (no actionable findings).

Checks:

  • Interruptibility: the helper-thread wait loops on concurrent.futures.wait(timeout=0.1), so a KeyboardInterrupt is raised promptly even where untimed lock waits are not interruptible. Verified with a real SIGINT (not the patched wait used by the tests) on Python 3.10.16 and 3.13.1: the stop request used the returned ID and the interrupt was re-raised.
  • No thread leak on the normal path: the daemon helper ends after setting the future (asserted by joining pyathena-spark-start threads in the token test). Without kill_on_interrupt, no thread is created (asserted).
  • A start failure without an interrupt still propagates as DatabaseError through future.result().
  • asyncio: without kill_on_interrupt, the start is awaited directly (no shield), so a cancelled task cannot leave an unretrieved task exception. With it, a second cancellation cancels the awaited start task and propagates.
  • Tests exercise the failing path: 28 of 33 new tests fail on the base code.

Recorded limitations (not introduced here): a KeyboardInterrupt raised inside Thread.start() itself, before the try, propagates without recovery as all interrupts during the start did before. The recovery wait is bounded by the connection's RetryConfig for the start request; a second interrupt ends it at once.

self,
session_id: str,
code_block: str,
description: str | None = None,
client_request_token: str | None = None,
) -> str:
"""Start a calculation execution with ``StartCalculationExecution``.

Without ``client_request_token``, a generated token is sent, so that a
retried request returns the calculation an earlier attempt started instead
of starting another one.

With ``kill_on_interrupt`` enabled, the request runs on a helper thread.
On ``KeyboardInterrupt``, the cursor first tries to abandon the request.
This succeeds only if the helper has not begun the request by then; the
helper then never sends it, and the interrupt propagates. Otherwise the
cursor waits for the request to finish, requests cancellation of the
calculation it started, waits for a terminal state, stores the calculation
ID and execution on the cursor, and re-raises the interrupt. Another
``KeyboardInterrupt`` during that wait propagates at once.

Args:
session_id: The session ID.
code_block: The code to run.
description: The calculation description.
client_request_token: The idempotency token of the request.

Returns:
The calculation execution ID.

Raises:
KeyboardInterrupt: If interrupted while starting the calculation. A
failure to start, cancel, or wait for the calculation becomes its
``__cause__``.
DatabaseError: If the request fails.
"""
request = self._build_start_calculation_execution_request(
session_id=session_id,
code_block=code_block,
description=description,
client_request_token=client_request_token or str(uuid.uuid4()),
)
if not self._kill_on_interrupt:
return self.__start_calculation(request)

future: Future[str] = Future()

def start() -> None:
# Begin the request only if no interrupt has given up on it yet.
if not future.set_running_or_notify_cancel():
return
try:
future.set_result(self.__start_calculation(request))
except BaseException as e:
future.set_exception(e)

try:
threading.Thread(target=start, name="pyathena-spark-start", daemon=True).start()

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent review (relayed) — Codex CLI 0.157.0, codex exec --sandbox read-only, model_reasoning_effort=high, run on a detached snapshot of c0188ced93137cd106a416bd092affaf13c4a552 against merge-base 0ba0875c72c702e0b893843d3c50aa98b4a61b8f, without the PR description or prior findings. Static review; the reviewer ran no tests.

Coverage reported: the full diff (Spark cursor changes, offline tests, docs/spark.md), the execute() callers, the request builder, and the retry helpers; confirmed pyathena/common.py and pyathena/aio/common.py are unchanged.

Verdict: FINDINGS (1)

  • P2 — pyathena/spark/common.py:407 (at c0188ce): the helper thread was started before the try that handles KeyboardInterrupt. A Ctrl-C while Thread.start() waits for the thread to initialize escapes without recovery, while the helper can still send StartCalculationExecution, leaving the calculation running with no ID on the cursor.

Verified and repaired in 6d0e5e9: Thread.start() moved inside the try; the helper claims the future with set_running_or_notify_cancel() before sending, and the recovery first calls future.cancel() — if that succeeds the request was never sent and never will be, so the interrupt propagates at once; otherwise the recovery waits for the request as before. New test test_calculate_interrupted_before_request_is_sent (helper not started / started after the interrupt) fails 2 of 4 cases on c0188ce and passes now.
The asyncio path has no equivalent window: _calculate has no await before asyncio.shield(start), and the start task's first step, which dispatches the request to the worker thread, is scheduled before the outer task can resume with CancelledError.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Repair self-review of 6d0e5e9 (range-diff 0ba0875..c0188ce → 0ba0875..6d0e5e9: the first two commits are unchanged; one new commit touching pyathena/spark/common.py and tests/pyathena/spark/test_common.py).

Round one (behavior): CLEAN. Future.set_running_or_notify_cancel() and Future.cancel() serialize on the future's condition lock, so exactly one outcome holds. Either the helper claims the future first, in which case cancel() returns False and the recovery waits for the ID as before, or cancel() wins and the helper returns without sending. An interrupt before _start_new_thread leaves a pending future that cancel() settles, so the recovery never waits on a thread that does not exist. The normal path and kill_on_interrupt=False are unchanged; __wait_for_start is never called on a cancelled future. Offline Spark tests: 83 passed. A real-SIGINT probe on 3.13.1 still stops the returned ID.

Round two (claims): CLEAN. The commit message matches the control flow. The _calculate docstring now states that an interrupt before the request is sent propagates without sending it. The docs/spark.md sentence ("first waits for the StartCalculationExecution request to finish") remains accurate because nothing is sent in the new case.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-up (relayed): Codex CLI 0.157.0, read-only, model_reasoning_effort=high, run on a detached snapshot of 6d0e5e9. It reviewed the repair diff c0188ce..6d0e5e9 and the range-diff, and did not see this thread. It was a static review; the reviewer ran no tests.

Verdict: FINDINGS (2 × P2). Both findings were checked against the code, and neither leaves a calculation running without a stop request:

  1. pyathena/spark/common.py:405 (at 6d0e5e9). An interrupt between the helper's set_running_or_notify_cancel() and the HTTP call makes cancel() return False, so the request is still sent and then stopped.
    Disposition: behavior kept, docstring corrected in a2008d2. The claim is the commit point: once the helper has begun, the cursor waits for the calculation ID and stops that calculation. That is the agreed design, "finish the in-flight start, then stop it." The docstring had said "before the helper sends the request", which over-promised. It now says "before the helper begins the request", and the two code comments were reworded to match.
  2. pyathena/aio/spark/cursor.py:154. A cancellation requested after ensure_future() but before the start task's first step still dispatches the request, which is then stopped.
    Disposition: rejected, no change. The outcome is the documented one: the cursor waits for the start, stops the calculation, and re-raises. The calculation is never left running unseen, which is what Spark: an interrupt during StartCalculationExecution leaves the calculation running #841 is about. Detecting "cancel requested but not yet delivered" before the first step would need Task.cancelling(), which is 3.11+ (the floor is 3.10), or private task state. AioSparkCursor._calculate documents that a cancellation waits for the request to finish, and that claim holds.

Repair self-review of a2008d2 (docstring and comments only):

  • Round one: CLEAN. No code change; lint passed.
  • Round two: CLEAN. The new wording matches the control flow in both outcomes of the set_running_or_notify_cancel() / cancel() race.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Independent follow-ups on the docstring repair (relayed). Codex CLI 0.157.0, read-only, model_reasoning_effort=high, detached snapshots, run without this thread. Static reviews only.

  • a2008d2, reviewing 6d0e5e9..a2008d2: FINDINGS, P2 at pyathena/spark/common.py:370. "An interrupt before the helper begins the request … does not send it" is false in one race: the interrupt can be raised while the future is pending, and the helper can then claim it before the handler calls cancel(). The outcome is decided by which call wins, not by when the signal arrives. Verified; reworded in 9a9965a.
  • 9a9965a, reviewing a2008d2..9a9965a: FINDINGS, P2 at the same line. "When handled, a request the helper has not yet begun is abandoned" still did not make the handler's cancel() attempt the condition. Verified; reworded in eeafbc1: "the cursor first tries to abandon the request. This succeeds only if the helper has not begun the request by then …"
  • eeafbc1, reviewing 9a9965a..eeafbc1: CLEAN.

Repair self-review of 9a9965a and eeafbc1: these are docstring-only changes.

  • Round one is CLEAN: there is no code change, and just lint passed.
  • Round two is CLEAN: the text now matches both outcomes of the race, and the code comments ("Begin the request only if no interrupt has given up on it yet", "The helper has not begun the request and never will") describe the successful-cancel() branch they sit in.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Rebase onto ceb9ad7 (#844 merged): repair review. The series moved from 0ba0875..eeafbc1 to ceb9ad7..bdfa630, and the PR went back to Draft for the review.

git range-diff:

Self-review round one (behavior): CLEAN.

  • Terminate only the Spark sessions a cursor started on close #844 changes __init__ and close() only. _calculate and __cancel_and_wait stop by calculation ID and never terminate a session, so a borrowed session or a per-call execute(session_id=...) override is unaffected.
  • just lint passed. pytest --noconftest over tests/pyathena/spark and tests/pyathena/aio/spark: 113 passed and 2 skipped; the 14 errors are live-fixture setups that --noconftest skips by design.

Self-review round two (claims): CLEAN.

Independent follow-up (relayed): CLEAN.

  • Reviewer: Codex CLI 0.157.0, read-only, model_reasoning_effort=high, run on a detached snapshot of bdfa630 with the range-diff and the full ceb9ad7..bdfa630 delta, without this thread. Static review; it ran no tests.
  • Codex reported that all six commits carried across, and that the conflict resolution kept both the ownership test and the start-cancellation tests.
  • It found that recovery cancels by calculation ID even under execute(session_id=...), that closing a borrowed cursor leaves its session running, and no actionable interaction.

return self.__wait_for_start(future)
except KeyboardInterrupt as interrupt:
if future.cancel():
# The helper has not begun the request and never will.
raise
_logger.warning("Query canceled by user.")
try:
self._calculation_id = self.__wait_for_start(future)
self.__cancel_and_wait(self._calculation_id)
except Exception as e:
raise interrupt from e
raise

def __start_calculation(self, request: dict[str, Any]) -> str:
"""Send a ``StartCalculationExecution`` request.

Args:
request: The request parameters.

Returns:
The calculation execution ID.

Raises:
DatabaseError: If the request fails.
"""
try:
response = retry_api_call(
self._connection.client.start_calculation_execution,
config=self._retry_config,
logger=_logger,
**request,
)
except Exception as e:
_logger.exception("Failed to execute calculation.")
raise DatabaseError(*e.args) from e
return cast(str, response.get("CalculationExecutionId"))

@staticmethod
def __wait_for_start(future: Future[str]) -> str:
"""Wait for the start request on a helper thread to finish.

Args:
future: The future of the start request.

Returns:
The calculation execution ID.

Raises:
DatabaseError: If the request failed.
"""
while not future.done():
wait((future,), timeout=_INTERRUPT_CHECK_INTERVAL)
return future.result()

def _cancel(self, query_id: str) -> None:
"""Stop a calculation execution with ``StopCalculationExecution``.

Expand Down
Loading
Loading