Cancel a Spark calculation interrupted while it is being started - #861
Conversation
| self._cancel(calculation_id) | ||
| self._calculation_execution = cast(AthenaCalculationExecution, self.__poll(calculation_id)) | ||
|
|
||
| def _calculate( |
There was a problem hiding this comment.
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 aKeyboardInterruptis raised promptly even where untimed lock waits are not interruptible. Verified with a realSIGINT(not the patchedwaitused 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-startthreads in the token test). Withoutkill_on_interrupt, no thread is created (asserted). - A start failure without an interrupt still propagates as
DatabaseErrorthroughfuture.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.
|
|
||
| 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: |
There was a problem hiding this comment.
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
ClientRequestTokenasidempotencyTokenforStartQueryExecutionbut not forStartCalculationExecution/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 withInvalidRequestException; 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:
- 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.
- Docs did not say that a caller-supplied
client_request_tokenreused 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.
| future.set_exception(e) | ||
|
|
||
| try: | ||
| threading.Thread(target=start, name="pyathena-spark-start", daemon=True).start() |
There was a problem hiding this comment.
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(atc0188ce): the helper thread was started before thetrythat handlesKeyboardInterrupt. A Ctrl-C whileThread.start()waits for the thread to initialize escapes without recovery, while the helper can still sendStartCalculationExecution, 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.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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:
pyathena/spark/common.py:405(at6d0e5e9). An interrupt between the helper'sset_running_or_notify_cancel()and the HTTP call makescancel()return False, so the request is still sent and then stopped.
Disposition: behavior kept, docstring corrected ina2008d2. 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.pyathena/aio/spark/cursor.py:154. A cancellation requested afterensure_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 needTask.cancelling(), which is 3.11+ (the floor is 3.10), or private task state.AioSparkCursor._calculatedocuments 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.
There was a problem hiding this comment.
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, reviewing6d0e5e9..a2008d2: FINDINGS, P2 atpyathena/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 callscancel(). The outcome is decided by which call wins, not by when the signal arrives. Verified; reworded in9a9965a.9a9965a, reviewinga2008d2..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'scancel()attempt the condition. Verified; reworded ineeafbc1: "the cursor first tries to abandon the request. This succeeds only if the helper has not begun the request by then …"eeafbc1, reviewing9a9965a..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 lintpassed. - 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.
eeafbc1 to
bdfa630
Compare
| future.set_exception(e) | ||
|
|
||
| try: | ||
| threading.Thread(target=start, name="pyathena-spark-start", daemon=True).start() |
There was a problem hiding this comment.
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:
- Commits 2–6 are identical.
- Commit 1 differs only in
tests/pyathena/aio/spark/test_cursor.py. That was an add/add conflict next to Terminate only the Spark sessions a cursor started on close #844'stest_session_ownership, resolved by keeping both sides: the imports are merged (AthenaSessionStatusfrom Terminate only the Spark sessions a cursor started on close #844, plusDatabaseErrorandRetryConfig), and the Terminate only the Spark sessions a cursor started on close #844 test is kept ahead of the new start-cancellation tests. pyathena/spark/common.py,pyathena/aio/spark/cursor.pyanddocs/spark.mdapplied cleanly.
Self-review round one (behavior): CLEAN.
- Terminate only the Spark sessions a cursor started on close #844 changes
__init__andclose()only._calculateand__cancel_and_waitstop by calculation ID and never terminate a session, so a borrowed session or a per-callexecute(session_id=...)override is unaffected. just lintpassed.pytest --noconftestovertests/pyathena/sparkandtests/pyathena/aio/spark: 113 passed and 2 skipped; the 14 errors are live-fixture setups that--noconftestskips by design.
Self-review round two (claims): CLEAN.
- The docs paragraphs sit after Terminate only the Spark sessions a cursor started on close #844's rewritten lifecycle section and do not contradict it.
just docs lintreports 0 errors.
Independent follow-up (relayed): CLEAN.
- Reviewer: Codex CLI 0.157.0, read-only,
model_reasoning_effort=high, run on a detached snapshot ofbdfa630with the range-diff and the fullceb9ad7..bdfa630delta, 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.
An interrupt that arrived while execute() was sending StartCalculationExecution left the calculation running with no ID on the cursor, and the session rejected new calculations until it finished. With kill_on_interrupt, the sync cursors now send the start request from a helper thread and the asyncio cursor shields it from cancellation; on an interrupt they wait for the request to finish, cancel the calculation it started, wait for a terminal state, and re-raise the interrupt, as the polling phase already does since #833. The Spark cursors also send a generated ClientRequestToken when the caller passes none. botocore does not generate one for StartCalculationExecution, so a retried start could fail on the busy session and leave the first calculation running unseen; with the token, the retry returns that calculation instead. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
An interrupt raised while the helper thread was being started escaped before the recovery, but the helper could still send the request and leave the calculation running unseen. The helper now claims the future before sending, and the recovery cancels an unclaimed future and re-raises at once, so a request is either never sent or waited for. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The interval exists so that the wait for the start request wakes up to check for Ctrl-C: before Python 3.14, an untimed lock wait on Windows cannot be interrupted by signals. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
bdfa630 to
1466e7b
Compare
| # 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 |
There was a problem hiding this comment.
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 of1466e7bwith 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.
Move the Spark start-phase interrupt handling from #861 into shared helpers and use them for the SQL cursors too. With kill_on_interrupt, StartQueryExecution now runs on a helper thread (sync) or is shielded from task cancellation (asyncio); an interrupt waits for the request, records the query ID on cursors that expose query_id, stops the query, waits for a terminal state, and re-raises. A request the helper has not begun is abandoned and never sent. The polling-phase recovery of the SQL and Spark cursors now goes through the same helpers, so the Spark cursors keep their behavior with less duplicated code. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
WHAT
SparkCursor, andAsyncSparkCursorwhoseexecute()starts the calculation on the caller's thread): withkill_on_interrupt,SparkBaseCursor._calculate()sendsStartCalculationExecutionfrom a short-lived helper thread and waits for it on the calling thread. OnKeyboardInterrupt, it waits for the request to finish, requestsStopCalculationExecutionfor the calculation it started, waits for a terminal state, storescalculation_idand the final execution, and re-raises the interrupt.AioSparkCursor): withkill_on_interrupt, the start request is wrapped in a task and awaited throughasyncio.shield(). Onasyncio.CancelledError, the same recovery runs and the cancellation is re-raised, sotask.cancelled()is true andasyncio.wait_for()raisesTimeoutError.__cause__. A second interrupt during the recovery propagates at once without cancelling.kill_on_interrupt=Falseis unchanged: the request runs on the calling thread / directly in the task, and the interrupt propagates immediately.ClientRequestToken(uuid4, one perexecute()) when the caller passes noclient_request_token; a caller-supplied token is sent as is. botocore does not auto-generate this token forStartCalculationExecution(it does forStartQueryExecution), so a retried start after a lost response is a new request; on the session made busy by the first attempt it fails withInvalidRequestException(measured for a different token on a busy session), leaving the first calculation running with no ID on the cursor.__cancel_and_wait()helper with the start phase.SparkBaseCursor(overridingBaseCursor._calculate) andAioSparkCursor;_build_start_calculation_execution_requestis reused as is.docs/spark.mdcancellation sections describe the start-phase behavior, the generated token, that a caller-supplied token must be unique per calculation, and that a stop request sent right after a calculation starts can occasionally have no effect.Behavior change (release note):
ClientRequestTokenunlessclient_request_tokenis given.kill_on_interrupt(default), an interrupt while the calculation is being started waits for the start request, then stops that calculation before the interrupt propagates. If the helper thread has not begun the request when the interrupt is handled, the request is abandoned and never sent.WHY
Closes #841. Parent: #791.
Measured with the CI account on
pyathena-spark(DPU 1/2/1, one session at a time, all sessions terminated afterwards):StartCalculationExecutionwith the same token, while the calculation runs or after it finished, returned the same calculation ID without a BUSY error and without a second calculation. A start with a different token on the busy session failed withInvalidRequestException.CodeBlockalso returned the original calculation without an error, so recovery by replaying a token could start code the user interrupted if the original request never arrived; recovery by listing calculations could stop another borrower's calculation on a shared session. This PR therefore waits for the in-flight request instead.CREATED) cancelled the calculation within 0.4 s in 5 of 6 attempts; in 1 of 6 it had no effect and the calculation endedCOMPLETEDafter about 95 s. Cancellation stays best-effort.TEST
Tested commits: 73a738e for the local live run and 6d0e5e9 for the offline tests (before rebasing). The later commits change only docstrings,
docs/spark.md, and the name and comment of the wait interval constant. The current head is 1466e7b (rebased onto 9a1a72d):just lintandjust docs lintpassed, and the offline Spark tests passed (113 passed, 2 skipped). The AWS suites run in CI.just format,just lint: passed.just docs lint: 0 errors.uv run sphinx-build -q -E docs <tmp>: no warnings fromdocs/spark.md.uv run --env-file .env pytest --noconftest -p no:xdist tests/pyathena/spark tests/pyathena/aio/spark -k "not (spark_dataframe or spark_sql or failed or test_cancel or executemany or context_manager)": 83 passed. The existing live aio testtest_context_manageralso ran against Athena with the generated token on 73a738e and passed.kill_on_interrupt; asyncio).KeyboardInterruptinjected into the wait whilestart_calculation_executionis blocked on an event; the request is sent once, the stop request uses the returned ID,calculation_idand the final execution are stored, and the same interrupt is re-raised (final stateCANCELEDandCOMPLETED). Start / cancel / wait failures become__cause__; a second interrupt propagates with the first as__context__and no stop; withoutkill_on_interruptno helper thread is created.task.cancel()while the start is blocked in the worker thread; the task stays pending until the request is released, then the stop is sent andtask.cancelled()is true.asyncio.wait_for()raisesTimeoutErrorafter the stop. Failures become__cause__; a second cancellation propagates without a stop; withoutkill_on_interruptthe cancellation propagates at once.kill_on_interrupt).uv run --env-file .env pytest -n 1 tests/pyathena/spark tests/pyathena/aio/spark→ 95 passed, 2 skipped, 286 s; afterwardsListSessionsshowed no non-terminated sessions. One unrelated Test run was in progress at the time.SIGINT(not a patched wait) delivered while the start request was blocked, on Python 3.10.16 and 3.13.1 (probe outside the repository, mocked client): the wait was interrupted, the stop request used the returned ID, andKeyboardInterruptwas re-raised.🤖 Generated with Claude Code