Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions docs/usage.md
Original file line number Diff line number Diff line change
Expand Up @@ -286,6 +286,7 @@ cursor.execute("SELECT * FROM one_row", cache_size=100, cache_expiration_time=36

Results will only be re-used from a succeeded DML query (the assumption being that you always want to re-run queries like `CREATE TABLE` and `DROP TABLE`)
whose query string (with `pyformat` parameters substituted) matches *exactly*, and that ran with the same schema and catalog as the cursor.
The cache is not used for a `qmark` query with parameters.
With `unload=True` on the pandas, Arrow, and Polars cursors, a query that is wrapped in `UNLOAD` is written to a new location each time, so the cache never matches it.

The S3 staging directory is not checked, so it's possible that the location of the results is not in your provided `s3_staging_dir`.
Expand Down
19 changes: 12 additions & 7 deletions pyathena/aio/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,8 @@ async def _execute( # type: ignore[override]
"""Start a query execution, or find a previous one to reuse.

The individual keyword arguments override the ``options`` field of the
same name unless None.
same name unless None. A query with execution parameters (``qmark``)
always starts a new execution.

Args:
operation: SQL query string.
Expand Down Expand Up @@ -85,12 +86,16 @@ async def _execute( # type: ignore[override]
paramstyle=paramstyle,
)
query, request = self._build_execute_request(operation, parameters, options)
query_id = await self._find_previous_query_id(
query,
options.work_group,
cache_size=options.cache_size,
cache_expiration_time=options.cache_expiration_time,
)
query_id = None
# Athena does not return the ExecutionParameters of earlier executions,
# so the cache cannot tell which parameters an execution ran with (#941).
if not request.get("ExecutionParameters"):
query_id = await self._find_previous_query_id(
query,
options.work_group,
cache_size=options.cache_size,
cache_expiration_time=options.cache_expiration_time,
)
if query_id is None:
query_id = await self._start_execution(lambda: self._start_query_execution(request))
return query_id
Expand Down
19 changes: 12 additions & 7 deletions pyathena/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -1280,7 +1280,8 @@ def _execute(
"""Start a query execution, or find a previous one to reuse.

The individual keyword arguments override the ``options`` field of the
same name unless None.
same name unless None. A query with execution parameters (``qmark``)
always starts a new execution.

Args:
operation: SQL query string.
Expand Down Expand Up @@ -1317,12 +1318,16 @@ def _execute(
paramstyle=paramstyle,
)
query, request = self._build_execute_request(operation, parameters, options)
query_id = self._find_previous_query_id(
query,
options.work_group,
cache_size=options.cache_size,
cache_expiration_time=options.cache_expiration_time,
)
query_id = None
# Athena does not return the ExecutionParameters of earlier executions,
# so the cache cannot tell which parameters an execution ran with (#941).

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 2 (claims, compatibility, operations) — FINDINGS (PR description only, corrected)

Scope: base e0e85da, head 9f3557c; claims in the PR body, commit messages, this comment, the _execute() / ExecuteOptions.cache_size docstrings, docs/usage.md:289, and the issue premise.

Claims checked:

  • Issue premise "AthenaQueryExecution.execution_parameters exposes them" — false in practice. Measured 2026-10-03 in work groups pyathena and primary: neither GetQueryExecution nor BatchGetQueryExecution (the API the cache search uses) returns ExecutionParameters for a query started with them, although botocore 1.43.102's QueryExecution shape declares the field. This comment's claim holds for both measured work groups.
  • Round 1's deferral: measured, SELECT ? AS v with paramstyle="qmark" and no parameters fails with INVALID_PARAMETER_USAGE: Incorrect number of parameters: expected 1 but found 0, so the remaining pre-existing match only affects invalid calls.
  • "Queries without parameters, including pyformat, keep the existing behavior": _build_start_query_execution_request() adds ExecutionParameters only for a non-empty list; pyformat passes None.
  • Docs reader: the parameterized-query section (docs/usage.md:108-143) says nothing about caching; the cache section now states the rule; no other doc mentions qmark with the cache.
  • Existing caller: no signature changes. Repeated qmark queries with identical parameters used to get correct cache hits and now run again (more Athena queries for those callers). This was missing from the release note.
  • AWS operator: the change only removes ListQueryExecutions / BatchGetQueryExecution calls; no retry layer changes.
  • Evidence: the local AWS run (21 passed) and offline revert check belong to 9126987; 9f3557c changes one docs sentence only.

Corrections to the PR description: named the second measured work group, added the identical-parameters behavior change to the release note, and stated which commit the test results belong to.

if not request.get("ExecutionParameters"):

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 1 (implementation behavior) — FINDINGS (1, repaired)

Scope: base e0e85da, reviewed head 9126987; all 6 changed files.

Covered:

  • Callers: every cursor (cursor.py, async_cursor.py, pandas/arrow/polars/s3fs sync + async, aio cursor/pandas/arrow/polars/s3fs) reaches the cache only through BaseCursor._execute() / AioBaseCursor._execute(), so the guard here and at pyathena/aio/common.py:92 covers all of them. Spark cursors do not use the cache.
  • Request shape: _build_start_query_execution_request() sets ExecutionParameters only when the list is non-empty, so qmark with None/[] and all pyformat queries keep the existing lookup; the global pyathena.paramstyle = "qmark" takes the same _prepare_query() path.
  • Behavior: return value, cursor state, on_start_query_execution, and interrupt handling are unchanged; a skipped lookup takes the same _start_execution() path as a cache miss. _find_previous_query_id() signatures are unchanged (dbt-athena 1.x does not call them; the legacy _execute() kwargs are untouched and the passthrough tests still assert the lookup call).
  • Tests: test_execute_qmark_parameters_skip_cache (sync/aio) fails with the pyathena/ change reverted; test_cache_size_with_qmark_parameters runs the issue's reproduction against Athena.

Finding: the added docs/usage.md:289 sentence explained why ("because Athena does not return ..."); user docs state behavior only, the rationale stays in the PR. Repaired in 9f3557c.

Deferred (pre-existing, not introduced): a qmark query text with ? placeholders executed later without parameters could still match an earlier parameterized execution of the same text and return its result instead of Athena's missing-parameter error. Detecting this would need SQL parsing; such a call is invalid input either way.

query_id = self._find_previous_query_id(
query,
options.work_group,
cache_size=options.cache_size,
cache_expiration_time=options.cache_expiration_time,
)
if query_id is None:
query_id = self._start_execution(lambda: self._start_query_execution(request))
return query_id
Expand Down
1 change: 1 addition & 0 deletions pyathena/options.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ class ExecuteOptions:
caching. 0 (default) disables the cache lookup, unless
``cache_expiration_time`` is set to a positive value, in which
case all queries within the expiration window are scanned.
A ``qmark`` query with parameters is never looked up.
cache_expiration_time: Maximum age in seconds of a cached query
result to consider for reuse. 0 (default) means no age limit.
result_reuse_enable: Enable Athena server-side result reuse for this
Expand Down
31 changes: 31 additions & 0 deletions tests/pyathena/aio/test_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -242,6 +242,37 @@ async def test_execute_internal_legacy_kwargs_passthrough(self):
cache_expiration_time=100,
)

async def test_execute_qmark_parameters_skip_cache(self):
"""A qmark query with parameters never searches the cache (no AWS, #941).

Mirrors the synchronous cursor test.
"""
cursor = AioCursor.__new__(AioCursor) # bypass __init__ to avoid AWS calls
cursor._connection = MagicMock()
cursor._connection.client.start_query_execution.return_value = {
"QueryExecutionId": "test_query_id"
}
cursor._retry_config = RetryConfig()
cursor._kill_on_interrupt = True

with (
patch.object(
AioCursor,
"_build_start_query_execution_request",
return_value={"ExecutionParameters": ["'1'"]},
) as request_mock,
patch.object(
AioCursor, "_find_previous_query_id", new_callable=AsyncMock, return_value="cached"
) as cache_mock,
):
query_id = await cursor._execute(
"SELECT ?", ["'1'"], paramstyle="qmark", cache_size=10, cache_expiration_time=100
)

assert query_id == "test_query_id"
assert request_mock.call_args.kwargs["execution_parameters"] == ["'1'"]
cache_mock.assert_not_awaited()

@pytest.mark.parametrize(
"final_state",
[AthenaQueryExecution.STATE_CANCELLED, AthenaQueryExecution.STATE_SUCCEEDED],
Expand Down
44 changes: 44 additions & 0 deletions tests/pyathena/test_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -199,6 +199,24 @@ def test_cache_size_with_work_group(self, cursor):
assert first_query_id != second_query_id
assert third_query_id in [first_query_id, second_query_id]

@pytest.mark.parametrize("cursor", [{"work_group": ENV.work_group}], indirect=["cursor"])
def test_cache_size_with_qmark_parameters(self, cursor):
query = f"SELECT ? AS v -- {datetime.now(UTC)!s}"

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) — FINDINGS (2; 1 rejected with evidence, 1 pre-existing and deferred)

Reviewer: Codex CLI 0.160.0, model gpt-6-astra (OpenAI), codex exec --sandbox read-only, session 01a10064-26ca-7751-a71d-935256265a41. Static review only: no builds, tests, edits, or network access. The prompt gave the base/head SHAs and the repository conventions, without the PR number, description, commit messages, or self-review records.
Snapshot: detached worktree at head 9f3557c, base e0e85da; afterwards it was still clean at the same HEAD, and the PR worktree was unchanged.

Covered (reviewer): the full diff; all sync, threaded-async, and aio SQL cursor paths (Dict, pandas, Arrow, Polars, S3FS); UNLOAD; legacy private-method signatures; cache matching; changed docs; relevant tests. The reviewer confirmed that the guard reaches every applicable cursor, private signatures are preserved, both mocked tests fail against the original code, and botocore's service documentation says the parameters are not returned.

  1. Introduced, P2 (this line): "SELECT ? AS v puts a parameter in the SELECT list, which Athena supports only in WHERE clauses, so the first execution fails." — Rejected. This test passed against Athena in the local run (-k "cache or qmark or legacy_kwargs": 21 passed) and returned [('2',)] / [('1',)]; the issue's live reproduction used the same SELECT ? AS v form.
  2. Pre-existing, P2 (pyathena/common.py:1125): the same ? SQL executed with parameters=None/[] and cache_size can reuse an earlier parameterized execution instead of letting Athena reject the missing binding. — Deferred, pre-existing (also present at the base SHA, as noted in round 1). Measured: Athena rejects such a call with INVALID_PARAMETER_USAGE: Incorrect number of parameters: expected 1 but found 0, so this only affects invalid calls. Skipping on unresolved ? markers would need a heuristic over the SQL text, which is outside this fix's scope.

No code change resulted, so no follow-up review is needed. The PR description now also quotes botocore's field documentation.


cursor.execute(query, ["'1'"], paramstyle="qmark")
first_query_id = cursor.query_id

# Different parameters must not reuse the earlier execution (#941).
cursor.execute(query, ["'2'"], paramstyle="qmark", cache_size=100)
assert cursor.query_id != first_query_id
assert cursor.fetchall() == [("2",)]

# Athena does not return the parameters of earlier executions,
# so even the same parameters run again.
cursor.execute(query, ["'1'"], paramstyle="qmark", cache_size=100)
assert cursor.query_id != first_query_id
assert cursor.fetchall() == [("1",)]

def test_cache_expiration_time(self, cursor):
query = f"SELECT * FROM one_row -- {datetime.now(UTC)!s}"

Expand Down Expand Up @@ -1303,6 +1321,32 @@ def test_execute_internal_legacy_kwargs_passthrough(self):
cache_expiration_time=100,
)

def test_execute_qmark_parameters_skip_cache(self):
"""A qmark query with parameters never searches the cache (no AWS, #941)."""
cursor = Cursor.__new__(Cursor) # bypass __init__ to avoid AWS calls
cursor._connection = MagicMock()
cursor._connection.client.start_query_execution.return_value = {
"QueryExecutionId": "test_query_id"
}
cursor._retry_config = RetryConfig()
cursor._kill_on_interrupt = True

with (
patch.object(
Cursor,
"_build_start_query_execution_request",
return_value={"ExecutionParameters": ["'1'"]},
) as request_mock,
patch.object(Cursor, "_find_previous_query_id", return_value="cached") as cache_mock,
):
query_id = cursor._execute(
"SELECT ?", ["'1'"], paramstyle="qmark", cache_size=10, cache_expiration_time=100
)

assert query_id == "test_query_id"
assert request_mock.call_args.kwargs["execution_parameters"] == ["'1'"]
cache_mock.assert_not_called()

def test_connection_level_callback(self):
"""Test connection-level default callback."""
callback_results = []
Expand Down
Loading