Re-raise and stop interrupted queries in the SQL cursors, sharing the Spark handling - #853
laughingman7743 wants to merge 6 commits into
Conversation
| return query_execution | ||
| self.__poll(query_id) | ||
| except Exception as e: | ||
| raise interrupt from e |
There was a problem hiding this comment.
Self-review round 1: implementation behavior. Result: CLEAN
Scope: c01c56f73c7dbf973fe73b52b6093f0a32952321..2de91e768196dc061251f7609ae0d9b9be2905cd (full diff: pyathena/common.py, pyathena/aio/common.py, both test files, docs/usage.md, docs/aio.md).
Covered:
- Callers of the shared
_poll():execute()ofCursor/DictCursor, pandas/arrow/polars/s3fs (sync), andAioCursorplus the aio pandas/arrow/polars/s3fs cursors. The thread-poolAsync*cursors call_poll()from executor threads (pyathena/async_cursor.py:153,pyathena/pandas/async_cursor.py:131, etc.), which never receiveKeyboardInterrupt, so they are unaffected. The Spark cursors override_poll(). SQLAlchemy does not call_poll()directly. - Cursor state after an interrupt:
_reset_state()already clearedresult_setandrowcount,query_idis set before polling, and no result set is built on the interrupt path, so nothing is left open. executemany(): bothpyathena/result_set.py:1064andpyathena/aio/common.py:651catchBaseException, close, and re-raise, so an interrupt stops the loop withrowcount == -1andquery_idkept. Before this change, an interrupted execution that endedSUCCEEDEDlet the loop continue with the next parameters.- Exception flow: the bare
raiseafter the innertry/except Exceptionre-raises the outer interrupt. A second interrupt or cancellation during the wait is aBaseException, so it propagates instead of becoming a cause (same as Define best-effort Spark calculation cancellation #833). - Tests: the new tests fail on the original
_poll()for the defect itself. TheSUCCEEDEDcases fail withDID NOT RAISE, because the fake result-set class lets the old code return normally.
No actionable findings. Round two will add the executemany() consequence to the PR description.
| assert cursor.query_id == "query_id" | ||
| assert cursor.result_set is None | ||
|
|
||
| async def test_execute_kill_on_interrupt_timeout(self): |
There was a problem hiding this comment.
Self-review round 1: test determinism
The first status request blocks on an event that is never set, so the 0.01 s timeout always fires while polling.
The re-poll after the stop request returns at once, so asyncio.wait_for() sees CancelledError and raises asyncio.TimeoutError on both the 3.10 wait_for implementation and the timeout()-based one from 3.12.
asyncio.TimeoutError is used instead of the builtin TimeoutError, because the two are distinct on 3.10.
|
|
||
| With `kill_on_interrupt` enabled, which is the default, a `KeyboardInterrupt` while `execute()` waits for the query | ||
| requests cancellation, waits until the query reaches a terminal state, and then propagates. | ||
| Cancellation is a best-effort request, so the query can still end as `SUCCEEDED` or `FAILED`. |
There was a problem hiding this comment.
Self-review round 2: claims, callers, and operations. Result: FINDINGS (PR description only, repaired)
Scope: c01c56f73c7dbf973fe73b52b6093f0a32952321..2de91e768196dc061251f7609ae0d9b9be2905cd, all claims in the PR description, commit message, _poll() docstrings, docs/usage.md, and docs/aio.md.
Claims checked:
- Best-effort stop can still end
SUCCEEDED: measured with oneSELECT 1on the CI account.StopQueryExecutionon an alreadySUCCEEDEDquery returns HTTP 200 and the state staysSUCCEEDED. So the race re-raises the interrupt without a cause, and this sentence and the_poll()docstrings hold. - "If the cancellation request fails, ... as its cause": covered by the
failure[cancel]tests on both bases. kill_on_interrupt=Falsekeeps the query running: no stop request is made (test_execute_without_kill_on_interrupt).- Timeouts: the commit message claim holds. On the original code,
asyncio.wait_for()surfacesOperationalError(the new timeout test fails this way), and CPython'sTimeout.__aexit__converts onlyCancelledErrorintoTimeoutError. - Existing callers: dbt-athena's current connection manager uses boto3 directly, and its legacy PyAthena cursor overrides
_poll()(dbt-athena/src/dbt/adapters/athena/connections_legacy.py:167), so it is unaffected. Thekill_on_interruptparameter descriptions inpyathena/connection.py:228and the cursor docstrings remain accurate.
Findings, repaired in the PR description:
- The WHY claimed that
TaskGroupcould not handle the cancellation. ATaskGroupstill raises itsExceptionGroupwhen a sibling fails, so this is narrowed to the verifiedasyncio.wait_for()effect. - The description omitted the
executemany()consequence: it now stops at the interrupted execution, where an interrupted execution that endedSUCCEEDEDused to let the loop continue. This is added, along with the dbt-athena compatibility note and the live stop measurement.
Deferred (pre-existing, out of scope): the thread-pool Async* cursor docstrings (e.g. pyathena/arrow/async_cursor.py:91) say kill_on_interrupt cancels on keyboard interrupt, but their polling runs in executor threads, which never receive KeyboardInterrupt. This PR does not change that path.
| if not self._kill_on_interrupt: | ||
| raise | ||
| _logger.warning("Query canceled by user.") | ||
| try: |
There was a problem hiding this comment.
Independent review (relayed): Codex CLI 0.157.0, model gpt-6-sol, reasoning effort high, session 01a0dcbd-22a4-7b63-b2a6-f2802f6e8c27. Result: FINDINGS
Scope: c01c56f73c7dbf973fe73b52b6093f0a32952321..2de91e768196dc061251f7609ae0d9b9be2905cd. The reviewer ran in a --sandbox read-only detached snapshot of the head with no .env. The prompt contained only the diff range, a file list, and the review questions; it had no PR number, description, commit message, or prior findings. Static review: the reviewer ran no tests. The snapshot and the PR worktree were unchanged afterwards.
Covered (reviewer's words): "synchronous and native asyncio execute() and executemany() paths through the default, dict, pandas, Arrow, Polars, and S3FS cursors; cursor state and exception chaining; the thread-backed async and Spark overrides; the new tests and both documentation examples."
[P2] Pre-existing, exposed by the new contract: "A second KeyboardInterrupt or task cancellation during the stop request or follow-up poll escapes the handler because both are BaseException subclasses, outside except Exception. With the query still running, execute() exits before observing a terminal state. ... The behavior predates the diff, while the new documentation states the wait without this qualification." (also pyathena/aio/common.py:155)
There was a problem hiding this comment.
Disposition: documented; behavior kept (pre-existing). Verified: a second KeyboardInterrupt or task cancellation raised during _cancel() or the follow-up poll is a BaseException, so except Exception does not catch it and it propagates with the first interrupt as its context. This predates the diff, matches the Spark contract from #833, and gives users a way to stop waiting on a query that does not stop. da16975 documents it: docs/usage.md says "A second KeyboardInterrupt during that wait propagates without waiting for the terminal state.", and docs/aio.md has the equivalent sentence for a repeated task cancellation.
There was a problem hiding this comment.
Independent follow-up review (relayed): Codex CLI, model gpt-6-sol, reasoning effort high, session 01a0dcc3-f3ef-7833-a415-86e914ba8fa2. Scope: 2de91e768196dc061251f7609ae0d9b9be2905cd..da169757662dda486e41349789ce7a77e67ddc91 (same base c01c56f73c7dbf973fe73b52b6093f0a32952321), read-only snapshot, static review.
Covered surfaces: the four files changed in 2de91e7..da16975, traced through pyathena/aio/cursor.py, pyathena/aio/common.py, pyathena/aio/util.py, pyathena/cursor.py, and pyathena/common.py. This was a static, read-only review; I did not run tests.
FINDINGS
- P2 — docs/aio.md:142, docs/usage.md:511: The new qualification covers a second interruption while polling for the terminal state, but says nothing about one during the stop request. Both handlers call
_cancelbefore that follow-up poll, and theirexcept Exceptionblocks do not catch a secondCancelledErrororKeyboardInterrupt. For example, a second task cancellation while the stop request is awaiting its worker thread exits the handler without confirming that the stop request ran or that the query became terminal. The documentation still implies the first interruption completes those steps unless interrupted “during that wait.”
Prior items
- Resolved. docs/aio.md:145 limits the timeout claim to polling and explains that cancellation during startup can leave
query_idunset while the worker-thread request continues. - Not an actual defect. In this test,
_executeis an immediately completingAsyncMock; the first status request then enters anEvent.wait()and suspends. On Python 3.10,wait_for()schedules the child task before its timeout can cancel it; on 3.11–3.14, it awaits the coroutine inside the timeout context. The event loop therefore reaches that blocked status request before delivering the 10 ms timeout. A separate barrier is unnecessary for this fixture. - Resolved. Both tests now return
RUNNINGbefore a terminal state and assert that the polling callback saw both states before the original interruption propagated. - Unresolved for the stop-request window; resolved for follow-up polling. The source behavior described in the finding remains possible, and the added wording qualifies only “that wait.”
There was a problem hiding this comment.
Fixed in f528d50. The repeated-interrupt note now covers the cancellation request as well as the follow-up wait, in both docs/usage.md and docs/aio.md, and says the query can keep running.
Independent follow-up review (relayed): Codex CLI, model gpt-6-sol, reasoning effort high, session 01a0dcc8-cb13-7712-8373-34cdcfc3f7e1. Scope: da169757662dda486e41349789ce7a77e67ddc91..f528d50bc242ccacb06631081a7f21bda31d70a2, read-only snapshot, static review.
Covered surfaces: docs/aio.md “Task cancellation” and docs/usage.md “Query cancellation on interrupt,” checked against the five named source files. The new wording covers a second interruption during both the cancellation request and the follow-up poll.
CLEAN. No actionable inaccuracies found in either section. This was a read-only source review; no tests were run.
| If the cancellation request fails, `asyncio.CancelledError` is raised with the error as its cause. | ||
| With `kill_on_interrupt=False`, `asyncio.CancelledError` is raised immediately and the query keeps running. | ||
|
|
||
| A timeout from `asyncio.wait_for()` therefore cancels the query and raises `asyncio.TimeoutError`: |
There was a problem hiding this comment.
Independent review (relayed, Codex gpt-6-sol): [P2] introduced
"The wait_for() example says a timeout cancels the query. If the timeout occurs while start_query_execution is running in a worker thread, execute() has not assigned query_id; cancellation cannot stop the request, and Athena may start a query after the task exits. The example can print None while that query keeps running."
There was a problem hiding this comment.
Disposition: fixed in da16975. Verified: aio StartQueryExecution runs through asyncio.to_thread() (pyathena/aio/util.py:42), and query_id is assigned only when _execute() returns. The docs now limit the cancellation to a timeout that expires while execute() waits for the query, and state that a timeout during the start leaves query_id as None while the pending start request can still start the query. execute() has no await between the query_id assignment and _poll(), so there is no third window. The example prints query_id as a value that may be None.
| kill_on_interrupt=True, final_state=AthenaQueryExecution.STATE_CANCELLED | ||
| ) | ||
| with pytest.raises(asyncio.TimeoutError): | ||
| await asyncio.wait_for(cursor.execute("SELECT 1"), timeout=0.01) |
There was a problem hiding this comment.
Independent review (relayed, Codex gpt-6-sol): [P2] introduced
"The 10 ms timeout test has no barrier confirming that the task reached polling. Under a scheduling delay, it can time out before a query ID is assigned; cancel is then never awaited and the assertion fails intermittently."
There was a problem hiding this comment.
Disposition: rejected with evidence. No await in execute() yields to the event loop before the first status request blocks. _execute is an AsyncMock, which completes without suspending, and _call_on_start_query_execution is synchronous, so query_id is always assigned before the loop can run the timeout callback. On 3.12+, wait_for() awaits the coroutine inside the calling task under timeouts.timeout(), so the timer can fire only at the first suspension, the blocked poll. On 3.10/3.11, wait_for() wraps the coroutine in a task whose first step call_soon places in _ready before _run_once moves any expired timer into _ready, so that step runs first even after a scheduling delay. The test passed on 3.13.1 and 3.10.16 locally.
| def test_execute_kill_on_interrupt(self, final_state): | ||
| """An interrupt cancels the query, waits for it, and is re-raised (no AWS).""" | ||
| cursor, cancel = _offline_cursor(kill_on_interrupt=True) | ||
| cursor._get_query_execution = MagicMock( |
There was a problem hiding this comment.
Independent review (relayed, Codex gpt-6-sol): [P3] introduced
"The follow-up status is immediately terminal in the new tests. Given a query that stays RUNNING after the stop request, an implementation that polls once and raises before terminal state could still satisfy these tests. They exercise the prior implementation's failing path, but do not assert the promised wait through a nonterminal state." (also tests/pyathena/aio/test_cursor.py:197)
There was a problem hiding this comment.
Disposition: fixed in da16975. Both test_execute_kill_on_interrupt tests now report RUNNING after the stop request, then the terminal state. They assert through the on_poll hook that the polled states were [RUNNING, final_state] before the interrupt or cancellation was re-raised. Mutation check: replacing the follow-up __poll() with a single _get_query_execution() call in both _poll() implementations makes all 4 of these tests fail. With the real implementation, 13 targeted tests pass on 3.13.1, and 11 interrupt tests pass on 3.10.16.
With kill_on_interrupt enabled, the shared _poll() of the SQL cursors requested StopQueryExecution after a KeyboardInterrupt or task cancellation, waited for a terminal state, and then returned the final execution instead of re-raising. Callers saw an OperationalError for a CANCELLED or FAILED query, or a normal return when the query SUCCEEDED first, and asyncio.wait_for()/asyncio.timeout() could not turn the cancellation into a timeout. Re-raise the original interrupt after the stop request and the wait, as #833 does for the Spark cursors, with a failure to stop or wait as its cause. Closes #840 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Limit the documented asyncio timeout behavior to a timeout while the query is being waited for, note that a repeated interrupt or task cancellation skips the wait, and make the tests assert that the interrupt is re-raised only after a non-terminal poll reaches a terminal state. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A second interrupt or task cancellation also escapes while the cancellation request is in progress, not only during the follow-up wait. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
f528d50 to
ef17dec
Compare
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>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
| _INTERRUPT_CHECK_INTERVAL = 0.1 | ||
|
|
||
|
|
||
| def _start_interruptibly( |
There was a problem hiding this comment.
Self-review round 1 (expanded scope: start phase + shared helpers): implementation behavior. Result: FINDINGS (1, repaired)
Scope: full diff 659676c07e2c09397b7cbc5740cc10bfe6fe41cb..d66c6994034894313f6b68f23f9f5d421ac13632 (rebased onto #861). Because the contract expanded, this is a full pass, not a range-diff.
Covered:
_execute()callers: only the cursors'execute(). That isCursor, the pandas/arrow/polars/s3fs sync cursors, theirAsync*thread-pool variants (start on the caller's thread), and the fiveAio*cursors. SQLAlchemy callsexecute()only. dbt-athena legacy calls_execute()and gets the start-phase handling.- Shared helpers: the sync helper keeps Cancel a Spark calculation interrupted while it is being started #861's abandon check (
Future.set_running_or_notify_cancel()/cancel()) and its interruptible wait. Name mangling inside thelambdaresolves per class (_BaseCursor__start_query_execution,_AioBaseCursor__...,_SparkBaseCursor__...). The bareraiseafter the innerexcept Exceptionre-raises the outer interrupt. - Spark:
_poll/_calculatenow pass their own poll/stop callables.__stop_started_calculationsetscalculation_idbefore the cancel, as before, so a cancel failure still leaves the ID on the cursor. All Cancel a Spark calculation interrupted while it is being started #861 tests pass with only their patch targets moved. - SQL state:
_set_interrupted_query_id()is a no-op onBaseCursor(AsyncCursorhas noquery_id) and setsquery_idinWithFetch/WithAsyncFetch. It is called before the cancel. On a start failure,query_idstaysNone. - Poll-phase aio tests: they keep a mocked
_execute, so the earlier determinism argument for the 10 mswait_fortest still holds. The start-phase tests use a separate helper with the real_execute().
Finding (repaired in d66c699): docs/aio.md said query_id is None only if the timeout expires before the query is started, but a start request that fails after the timeout also leaves it None. Reworded to "only if no query was started".
Out of scope, same as the Spark cursors since #861: in asyncio, a start task that has not begun when the cancellation arrives still sends its request, which is then stopped. The sync helper abandons it.
| The `query_id` property returns that query's ID. | ||
| If the request has not been sent yet when the interrupt is handled, it is never sent. | ||
| `AsyncCursor` and its variants also stop a query whose start is interrupted in `execute()`. | ||
| They wait for queries on worker threads, which do not receive `KeyboardInterrupt`. |
There was a problem hiding this comment.
Self-review round 2 (expanded scope): claims, callers, and operations. Result: CLEAN
Scope: 659676c07e2c09397b7cbc5740cc10bfe6fe41cb..d66c6994034894313f6b68f23f9f5d421ac13632, the rewritten PR description, docstrings, and docs/usage.md / docs/aio.md.
Claims checked:
- No SQL token generation needed: botocore's Athena model marks
StartQueryExecution.ClientRequestTokenasidempotencyToken(auto-generated per call) andStartCalculationExecution's not.RetryConfigdefaults toTHROTTLING_ERROR_CODES, which do not start a query. - "
AsyncCursorand its variants also stop a query whose start is interrupted":pyathena/async_cursor.py:225and the pandas/arrow/polars/s3fs async cursors callself._execute()on the caller's thread. Polling runs inself._executor. - "If the request has not been sent yet ..., it is never sent": this is the sync abandon path, covered by the Cancel a Spark calculation interrupted while it is being started #861 tests
test_calculate_interrupted_before_request_is_sent[False/True], which now exercise the shared helper. - Real signal: a
SIGINTviaos.killwhileCursor.execute()was blocked in a mockedStartQueryExecutionproduced one start request, a stop with the returned ID,query_idset, andKeyboardInterruptwithout a cause. Checked on 3.13.1 and 3.10.16. - Operational: one short-lived daemon thread per
StartQueryExecutionwithkill_on_interrupt, and no extra AWS calls on the normal path. The start still goes throughretry_api_callwith the same config. - Evidence scope: local results are offline only (3.13.1 and 3.10.16). The AWS suites were not run locally for this revision. An accidental local run of all of
tests/pyathenawith--noconftestsent some queries from tests that callconnect()directly; its missing-schema failures are not used as evidence, and the PR description says so.
No corrections were needed beyond the round 1 repair.
Clear the previous calculation when a Spark cursor starts a new one, so that a failed execution no longer leaves the earlier calculation's state next to the new calculation ID. Make the asyncio start-phase timeout tests wait for the timeout before releasing the start request, and state when an asyncio timeout leaves query_id unset without implying that no query started. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
| async def test_execute_timeout_while_starting(self): | ||
| """A timeout during the start request stops the query and raises TimeoutError (no AWS).""" | ||
| cursor, cancel, started, release = _starting_cursor() | ||
| timer = threading.Timer(0.2, release.set) |
There was a problem hiding this comment.
Independent review (relayed): Codex CLI 0.157.0, reported model gpt-6-astra, reasoning effort high, session 01a0e6c8-41ad-7423-a641-afcd8c29f7fc. Result: FINDINGS
Scope: 659676c07e2c09397b7cbc5740cc10bfe6fe41cb..d66c6994034894313f6b68f23f9f5d421ac13632 (full expanded scope). The reviewer ran in a --sandbox read-only detached snapshot with no .env. The prompt contained no PR number, description, commit message, or prior findings. Static review. The snapshot and the PR worktree were unchanged afterwards.
Covered (reviewer's words): "the full diff and the sync, thread-pool async, and native asyncio cursor families, including pandas/Arrow/Polars/S3FS and Spark. Traced start/poll/cancel handling, exception identity/chaining, repeated interruption, helper lifecycle, cursor state, legacy methods and overrides, Python ≥3.10 compatibility, tests, and documentation." It found "No additional defect ... in the shared-helper refactor or legacy method signatures."
1. [P2] introduced: "The release timer starts before wait_for() establishes its timeout. If the test thread is descheduled for over 200 ms after timer.start(), the request is released before execution begins. The mocked query then reaches CANCELLED and raises OperationalError, failing the expected TimeoutError assertion despite correct implementation."
There was a problem hiding this comment.
Disposition: fixed in 10f7a83, together with the same pattern in the Spark test from #861 (tests/pyathena/aio/spark/test_cursor.py). Both tests now start wait_for() as a task, wait for the start request to begin, and then await asyncio.sleep(0.1) before releasing the request. The sleep begins after wait_for() scheduled its 0.05 s deadline, so its timer is due later. The event loop therefore runs the timeout callback first, and the cancellation reaches the task before the test resumes and releases the request. threading.Timer is gone. Both tests passed 30 of 30 repeated runs.
| With `kill_on_interrupt=False`, `asyncio.CancelledError` is raised immediately and the query keeps running. | ||
|
|
||
| A timeout from `asyncio.wait_for()` therefore cancels the query and raises `asyncio.TimeoutError`. | ||
| `query_id` is `None` only if no query was started, for example when the timeout expires while looking up a cached result. |
There was a problem hiding this comment.
Independent review (relayed, Codex gpt-6-astra): 2. [P3] introduced
"With kill_on_interrupt=False, a timeout during StartQueryExecution cancels the await while its underlying thread continues. Athena can start the query, but the cursor retains query_id=None. Repeated cancellation during the protected start wait can also abandon ID recovery. Document None as an unavailable ID, rather than evidence that no query exists."
There was a problem hiding this comment.
Disposition: fixed in 10f7a83. Verified: with kill_on_interrupt=False, or after a second cancellation during the shielded start, a query can start while query_id stays None. The sentence is now one-directional: "query_id is None if the timeout expires before the start request is sent, for example while looking up a cached result." The preceding paragraph already says that with kill_on_interrupt=False the query keeps running, and that another cancellation during the waits leaves it running.
| ) | ||
|
|
||
| future: Future[str] = Future() | ||
| def __stop_started_calculation(self, calculation_id: str) -> None: |
There was a problem hiding this comment.
Independent review (relayed, Codex gpt-6-astra): 3. [P2] pre-existing, preserved by the Spark refactor
Anchored here because the lines are outside the diff: pyathena/spark/cursor.py:147, pyathena/aio/spark/cursor.py:351.
"Execute calculation A successfully, then start B on the same cursor. Interrupt/cancel B while polling and make _cancel() fail. The exception propagates, but calculation_id identifies B while calculation_execution, state, and output accessors still describe A. Neither execution path clears the previous calculation object."
There was a problem hiding this comment.
Disposition: fixed in 10f7a83 (pre-existing; folded in as a contained Spark consistency fix). Verified: SparkCursor.execute() and AioSparkCursor.execute() overwrote _calculation_id without clearing _calculation_execution. Both now set them to None before _calculate(), like the SQL cursors' _reset_state(). The #833 poll-phase failure tests (sync and asyncio) now start with a previous calculation on the cursor and assert that calculation_id is the new one and calculation_execution is None. Without the reset, all 4 fail. Side effect, checked: cancel() during a new start now raises ProgrammingError instead of stopping the previous, finished calculation. The live test_cancel tests wait for an ID that is neither None nor the previous one, so they are unaffected. This is recorded as a release-note item in the PR description.
There was a problem hiding this comment.
Independent follow-up review (relayed): Codex CLI, reported model gpt-6-astra, reasoning effort high, session 01a0e710-57c6-7d73-b6b3-e9bb82e625e4. Scope: d66c6994034894313f6b68f23f9f5d421ac13632..10f7a832ebbac10440706167d1e370f90b23c2ab (same base 659676c07e2c09397b7cbc5740cc10bfe6fe41cb), read-only snapshot, static review. The snapshot and the PR worktree were unchanged afterwards.
Covered surfaces: all six changed files; Spark start/poll/cancellation paths; calculation_id, calculation_execution, state, and output accessors; repository callers; documentation; and local CPython asyncio sources for 3.10–3.14.
CLEAN — no actionable defects found in d66c699..10f7a832.
-
Timeout-test race: resolved. Both rewritten tests wait for the start request to begin, then release it only after the event-loop sleep.
- Python 3.10–3.11:
wait_for()registers its timeout before execution starts. Its deadline precedes the subsequently registered sleep deadline. Even if both timers become overdue together, the timeout callback queues thewait_for()continuation first; that continuation requests cancellation before the test resumes and releases the worker. - Python 3.12–3.14:
wait_for()uses the timeout context manager, whose earlier timer directly cancels the executing task before the sleep continuation releases the worker.
The pending-task assertion, subsequent
TimeoutError, cancellation-call assertion, and retained-ID assertion collectively exercise the intended behavior. The existing 10-second worker watchdog remains a finite scheduling limit; the original independent-thread release race is removed. - Python 3.10–3.11:
-
query_iddocumentation: resolved.docs/aio.md:148removes the exclusivity claim. It gives the pre-start/cache-lookup case without asserting thatNoneproves no query started, allowing the documented disabled-cleanup and repeated-cancellation cases. -
Stale Spark calculation: resolved.
pyathena/spark/cursor.py:148andpyathena/aio/spark/cursor.py:352clear both fields before starting another calculation. Start failures leave no previous ID/result; failures after obtaining the new ID retain that ID without the previous execution. Successful interruption cleanup still stores the current terminal execution.Compatibility is consistent:
stateand execution-derived properties already supportNone;cancel()rejects an unknown ID and targets the current calculation once known. The strengthened tests seed previous state and verify its removal for both cancellation-request and cleanup-wait failures. They directly protect the stale-execution regression, though they do not independently test clearing a previous ID when startup fails.
Static review only: no tests, builds, writes, or network access. HEAD remained 10f7a832; the worktree remained clean.
Author note: the remark that no test separately covers clearing a previous ID when the start fails is not taken up. The reset runs unconditionally before _calculate() (pyathena/spark/cursor.py:148, pyathena/aio/spark/cursor.py:352), and the new assertions already fail without it.
WHAT
With
kill_on_interruptenabled (the default), SQL cursors now handle an interrupt the way the Spark cursors do since #833 and #861, and the SQL and Spark cursors share one implementation of that handling.BaseCursor._poll(),AioBaseCursor._poll()): after the interrupt, the cursor requestsStopQueryExecutionand waits for a terminal state as before. It then re-raises the originalKeyboardInterrupt/asyncio.CancelledError, whatever the terminal state is. Before, the cursor returned the final execution.__cause__.query_idkeeps the query's ID.BaseCursor._execute(),AioBaseCursor._execute()), new for SQL:StartQueryExecutionruns on a short-lived helper thread (sync) or is shielded from task cancellation (asyncio).query_id, stops the query, waits for a terminal state, and re-raises._start_interruptibly()/_poll_interruptibly()inpyathena/common.py, and their asyncio counterparts inpyathena/aio/common.py, hold the helper-thread /asyncio.shield()logic that Cancel a Spark calculation interrupted while it is being started #861 added to the Spark cursors.SparkBaseCursorandAioSparkCursornow call these helpers with their own start / poll / stop callables.pyathena.common/pyathena.aio.commonloggers._set_interrupted_query_id()is a no-op hook onBaseCursor, overridden inWithFetch/WithAsyncFetchto setquery_id.Cursor,DictCursor, the pandas/arrow/polars/s3fs cursors, and theAio*cursors.Async*cursors, whoseexecute()starts the query on the caller's thread. They poll on worker threads, which never receiveKeyboardInterrupt.kill_on_interrupt=False: unchanged. The request runs on the caller's thread or directly in the task, and the interrupt propagates immediately.ClientRequestToken, unlike the Spark path in Cancel a Spark calculation interrupted while it is being started #861. botocore auto-generates it forStartQueryExecution(idempotencyTokenin the service model; not forStartCalculationExecution), and PyAthena's own retries default to throttling errors, which do not start a query.SparkCursor.execute()andAioSparkCursor.execute()now clearcalculation_idand the previous calculation execution before starting. Before, a failed execution (for example, a failed cancel request after an interrupt) left the earlier calculation'sstateand outputs next to the newcalculation_id.executemany(): it stops at the interrupted execution and re-raises, withrowcount == -1andquery_idkept. Before, an interrupted execution that endedSUCCEEDEDlet the loop continue._execute(), so it gets the start-phase handling. It overrides_poll(), so its polling behavior is unchanged. Its current connection manager uses boto3 directly.docs/usage.mdand "Task cancellation" indocs/aio.md.Behavior change (release note):
kill_on_interrupt=True, an interrupt duringexecute()now propagates asKeyboardInterrupt/asyncio.CancelledError. Before, it surfaced asOperationalError(query endedCANCELLED/FAILED) or was lost (query endedSUCCEEDED). Callers that caughtOperationalErrorafter Ctrl-C or task cancellation need to handle the interrupt instead.asyncio.wait_for()now raisesTimeoutError.StartQueryExecutionand stops the query it started.kill_on_interrupt=True, eachStartQueryExecutionof a synchronous cursor runs on a short-lived daemon thread.execute()onSparkCursor/AioSparkCursor,stateand the other calculation properties returnNone(or the new calculation's values) instead of the previous calculation's.WHY
Closes #840. Swallowing the interrupt lost Ctrl-C when the query finished first. In asyncio,
task.cancel()did not end with a cancelled task, and a timeout fromasyncio.wait_for()surfaced as the query'sOperationalError, or as a normal return, instead ofTimeoutError.After #861 fixed the start phase for Spark (#841), the maintainer asked this PR to cover the SQL start phase as well and to share the implementation with the Spark cursors.
TEST
Tested commit: 10f7a83.
just format,just lint: passed (ruff, format check, mypy, cfn-lint, license headers).just docs lint: 0 errors.uv run --env-file .env pytest --noconftest -p no:xdist ....--noconftestskips the session setup that creates AWS resources.tests/pyathena/test_cursor.py,tests/pyathena/aio/test_cursor.py,-k "interrupt or starting or on_poll_invoked or on_poll_none or legacy_kwargs_passthrough"): 26 passed on Python 3.13.1 and 3.10.16.tests/pyathena/spark tests/pyathena/aio/spark, excluding tests that need AWS fixtures): 105 passed on 3.13.1. On 3.10.16, a combined SQL and Spark selection passed (138 tests). The Cancel a Spark calculation interrupted while it is being started #861 tests pass unchanged except that the patch targets moved topyathena.commonand the wait-interrupting helper moved totests/pyathena/util.py.StartQueryExecutionis blocked. The request is sent once, the stop uses the returned ID, and the statesRUNNING→ terminal are polled before the same interrupt is re-raised.query_idis set, and no result set is built.__cause__. Withoutkill_on_interrupt, no helper thread is created.AsyncCursor.execute()also stops the query.task.cancel()while the start is blocked keeps the task pending until the request finishes, then stops the query, andtask.cancelled()is true.asyncio.wait_for()raisesTimeoutErrorafter the stop. Failures become__cause__. Withoutkill_on_interrupt, the cancellation propagates at once.asyncio.sleep()that ends after the timeout. They passed 30 of 30 repeated runs.kill_on_interrupt=Falsepath.SIGINT(probe outside the repository, mocked client) whileCursor.execute()was blocked inStartQueryExecution, on 3.13.1 and 3.10.16: one start request, stop with the returned ID,query_idset, andKeyboardInterruptre-raised without a cause.StopQueryExecutioncheck from the earlier revision of this PR still apply: a stop on aSUCCEEDEDquery returns HTTP 200 and the state staysSUCCEEDED.tests/pyathenawith--noconftestalso sent some queries to Athena from tests that callconnect()directly. Its failures were missing-schema errors caused by the skipped setup (two checked). It is not used as evidence here. The AWS suites run in CI once the PR is Ready.🤖 Generated with Claude Code