Skip to content
Closed
18 changes: 14 additions & 4 deletions .github/workflows/release.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -5,13 +5,23 @@ on:
tags:
- 'v*'

permissions:
id-token: write
contents: write
permissions: {}

jobs:
# Runs every suite on every supported Python version for the tagged commit;
# nothing is built or published unless all of them pass.
test:
uses: ./.github/workflows/test.yaml
permissions:
contents: read
id-token: write

release:
needs: test
runs-on: ubuntu-latest
permissions:
id-token: write
contents: write

env:
PYTHON_VERSION: '3.12'
Expand All @@ -22,7 +32,7 @@ jobs:

- uses: astral-sh/setup-uv@37802adc94f370d6bfd71619e3f0bf239e1f3b78 # v7.6.0
with:
python-version: ${{ matrix.python-version }}
python-version: ${{ env.PYTHON_VERSION }}
enable-cache: true

- name: Build
Expand Down
15 changes: 10 additions & 5 deletions .github/workflows/test-suite.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,15 @@ on:
test-type:
required: true
type: string
python-versions:
description: JSON array of the Python versions to test
required: true
type: string

jobs:
run:
# External fork contributions must validate AWS behavior in their own account.
if: github.event_name != 'pull_request' || github.event.pull_request.head.repo.full_name == github.repository
runs-on: ubuntu-latest

env:
Expand All @@ -18,16 +24,15 @@ jobs:
AWS_ATHENA_WORKGROUP: pyathena
AWS_ATHENA_SPARK_WORKGROUP: pyathena-spark
AWS_ATHENA_MANAGED_WORKGROUP: pyathena-managed
# Registered S3 Tables catalog (s3tablescatalog/<table-bucket>) and namespace
# for the SQLAlchemy S3 Tables tests; the table bucket, namespace, and the
# AWS analytics-services integration are provisioned out of band.
# The SQLAlchemy S3 Tables tests need a fixed namespace, which the test
# account no longer provides (master creates one per test session), so
# they are skipped on this branch.
AWS_ATHENA_S3_TABLES_CATALOG: s3tablescatalog/laughingman7743-pyathena-s3-tables
AWS_ATHENA_S3_TABLES_NAMESPACE: pyathena

strategy:
fail-fast: false
matrix:
python-version: ['3.10', '3.11', '3.12', '3.13', '3.14']
python-version: ${{ fromJSON(inputs.python-versions) }}

steps:
- name: Checkout
Expand Down
86 changes: 79 additions & 7 deletions .github/workflows/test.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -2,15 +2,23 @@ name: Test

on:
pull_request:
# ready_for_review starts the AWS jobs for a pull request leaving Draft;
# converted_to_draft starts a run without them, which cancels an
# in-progress run through the concurrency group.
types: [opened, synchronize, reopened, ready_for_review, converted_to_draft]
paths-ignore:
- 'docs/**'
- '**.md'
schedule:
- cron: '0 0 * * 0'
# Allows refreshing the README status badge on demand: the badge reflects
# the latest run on the default branch, which is otherwise only the weekly
# scheduled run and stays red for up to a week after a transient failure.
# Runs every suite on the selected branch, on the requested Python versions.
workflow_dispatch:
inputs:
python-versions:
description: Comma-separated Python versions, such as 3.12 or 3.11,3.14; empty for every supported version
type: string
default: ''
# The Release workflow runs every suite on every supported Python version
# before publishing.
workflow_call:

permissions:
id-token: write
Expand All @@ -23,19 +31,83 @@ concurrency:
cancel-in-progress: true

jobs:
# Offline checks run for every event, including Draft and fork pull requests.
lint:
runs-on: ubuntu-latest
permissions:
contents: read
steps:
- name: Checkout
uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2
with:
persist-credentials: false
- uses: astral-sh/setup-uv@37802adc94f370d6bfd71619e3f0bf239e1f3b78 # v7.6.0
with:
python-version: '3.12'
enable-cache: true
- uses: taiki-e/install-action@7a79fe8c3a13344501c80d99cae481c1c9085912 # v2.81.10
with:
tool: just
- run: just lint

# Selects the Python versions of the AWS suites. Draft and external-fork pull
# requests run none. A ready pull request tests the newest Python version; a
# dispatch tests the requested versions or every version, and the Release
# workflow every version.
versions:

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, and operations), by the authoring model; not an independent review.

Scope: base 4ef4d3e, head 610d6c3, covering the PR body, the commit messages, the changed docstrings and comments, and docs/sqlalchemy.md.

Result: FINDINGS. There were two gaps in the PR description; both are corrected, and the code is unchanged.

Claims checked:

  • The fixed S3 Tables namespace no longer exists. s3tables:GetNamespace(pyathena) returns NotFound. All three S3 Tables tests are gated on the namespace, so they skip.
  • Draft and fork PRs run no AWS suites. For Draft, this was observed on this PR's run: only lint ran. For forks, it follows from the conditions on versions and on test-suite.yaml's run job.
  • A tag runs every version before publishing. This holds in the code: release needs test, and the called workflow's event_name is push, so it selects the full list. It is not exercised until a tag is pushed, and neither is master's Run the full test matrix in the Release workflow before publishing #862.
  • "Scheduled runs only use the default branch's workflow" is GitHub's documented behavior, so the 3.x schedule never ran.
  • The reorder changed only the workflows: git diff 7c1db8f 610d6c3 touches only .github/workflows/*, so the local test results still apply.
  • STRUCT DDL, Double CAST, and INTEGER being accepted in DDL: covered by the round-one evidence and the Athena probe on fix: render Hive STRUCT syntax in table column DDL #870.

Existing callers:

  • 3.x allows sqlalchemy>=1.0.0. The compiler code only uses Column and a hasattr(types, "Double") guard, and the new Double test skips on SQLAlchemy versions below 2.0.
  • No other docs describe DDL with ROW(...). docs/usage.md:596 is a SELECT with CAST.

Operations:

  • A ready 3.x PR now runs the three suites on Python 3.14 only.
  • A tag runs the three suites on five Python versions, with 3.x's retry behavior and without master's rerun-once (Rerun tests once on Athena service-side query failures #811). A transient failure blocks publishing until the failed jobs are re-run.

Corrections: the PR body's TEST section now records the Draft CI observation, the namespace check, and the fact that the release gate has not run yet.

Separate finding, not in this PR: the test table bucket holds 539 namespaces, mostly pyathena_test_*. Per-session namespaces from master (#815) are accumulating. To be reported separately.

if: >-
github.event_name != 'pull_request' ||
(!github.event.pull_request.draft &&
github.event.pull_request.head.repo.full_name == github.repository)
runs-on: ubuntu-latest
permissions: {}
outputs:
python-versions: ${{ steps.select.outputs.python-versions }}
steps:
- id: select
env:
EVENT_NAME: ${{ github.event_name }}
REQUESTED_VERSIONS: ${{ inputs.python-versions }}
# Every supported version, oldest first; keep in sync with the
# pyproject.toml classifiers.
PYTHON_VERSIONS: '["3.10", "3.11", "3.12", "3.13", "3.14"]'
run: |
case "$EVENT_NAME" in
pull_request)
versions=$(jq -c '[last]' <<< "$PYTHON_VERSIONS")
;;
workflow_dispatch)
versions=$(jq -c --arg requested "$REQUESTED_VERSIONS" '
($requested | split(",") | map(gsub("\\s"; "")) | map(select(. != "")) | unique) as $selected
| if $selected == [] then .
elif ($selected - .) == [] then $selected
else error("unsupported Python versions: \($selected - . | join(", "))")
end' <<< "$PYTHON_VERSIONS")
;;
*)
# The Release workflow (a workflow_call from a tag push).
versions=$(jq -c '.' <<< "$PYTHON_VERSIONS")
;;
esac
echo "python-versions=$versions" >> "$GITHUB_OUTPUT"

test:
needs: versions
uses: ./.github/workflows/test-suite.yaml
with:
test-type: pyathena
python-versions: ${{ needs.versions.outputs.python-versions }}

test-sqla:
needs: [test]
needs: [versions, test]
uses: ./.github/workflows/test-suite.yaml
with:
test-type: sqla
python-versions: ${{ needs.versions.outputs.python-versions }}

test-sqla-async:
needs: [test-sqla]
needs: [versions, test-sqla]
uses: ./.github/workflows/test-suite.yaml
with:
test-type: sqla_async
python-versions: ${{ needs.versions.outputs.python-versions }}
8 changes: 5 additions & 3 deletions docs/sqlalchemy.md
Original file line number Diff line number Diff line change
Expand Up @@ -697,12 +697,14 @@ This generates the following SQL structure:

```sql
CREATE TABLE users (
id INTEGER,
profile ROW(name STRING, age INTEGER, email STRING),
settings ROW(theme STRING, notifications ROW(email STRING, push STRING))
id INT,
profile STRUCT<name:STRING, age:INTEGER, email:STRING>,
settings STRUCT<theme:STRING, notifications:STRUCT<email:STRING, push:STRING>>
)
```

`CREATE TABLE` renders `AthenaStruct` columns with Hive `STRUCT<name:type, ...>` syntax at every nesting depth, including STRUCT values inside MAP and ARRAY.

#### Querying STRUCT data

PyAthena automatically converts STRUCT data between different formats:
Expand Down
8 changes: 6 additions & 2 deletions pyathena/pandas/util.py
Original file line number Diff line number Diff line change
Expand Up @@ -261,13 +261,17 @@ def to_sql(
).Bucket(bucket_name)
cursor = conn.cursor()

# Athena stores identifiers in lowercase and information_schema reports them
# that way, so compare lowercase literals.
schema_literal = schema.lower().replace("'", "''")
name_literal = name.lower().replace("'", "''")
table = cursor.execute(
textwrap.dedent(
f"""
SELECT table_name
FROM information_schema.tables
WHERE table_schema = '{schema}'
AND table_name = '{name}'
WHERE table_schema = '{schema_literal}'
AND table_name = '{name_literal}'
"""
)
).fetchall()
Expand Down
41 changes: 37 additions & 4 deletions pyathena/spark/async_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,27 @@ def __init__(
max_workers: int = (cpu_count() or 1) * 5,
**kwargs,
):
"""Initialize the cursor and start or attach to a Spark session.

Args:
session_id: ID of an existing session to use. If omitted, a new
session is started.
description: Description of a new session.
engine_configuration: Engine configuration of a new session.
notebook_version: Notebook version of a new session.
session_idle_timeout_minutes: Idle timeout of a new session in minutes.
max_workers: Maximum number of threads for asynchronous operations.
**kwargs: Arguments passed to ``SparkBaseCursor``.

Raises:
ValueError: If ``max_workers`` is not greater than 0.
OperationalError: If the supplied session does not exist, or the
session cannot be started or does not become idle.
"""
# Created before the session so that an invalid max_workers cannot leave
# a newly started session behind; the executor starts no threads until used.
self._max_workers = max_workers
self._executor = ThreadPoolExecutor(max_workers=max_workers)
super().__init__(
session_id=session_id,
description=description,
Expand All @@ -83,12 +104,24 @@ def __init__(
session_idle_timeout_minutes=session_idle_timeout_minutes,
**kwargs,
)
self._max_workers = max_workers
self._executor = ThreadPoolExecutor(max_workers=max_workers)

def close(self, wait: bool = False) -> None:
super().close()
self._executor.shutdown(wait=wait)
"""Terminate the Spark session, then shut down the executor.

The executor is shut down even if terminating the session fails.
If termination fails, calling this method again retries it.

Args:
wait: Whether to wait for submitted futures to finish before returning
or raising.

Raises:
OperationalError: If terminating the session fails.
"""
try:
super().close()
finally:
self._executor.shutdown(wait=wait)

def calculation_execution(self, query_id: str) -> "Future[AthenaCalculationExecution]":
return self._executor.submit(self._get_calculation_execution, query_id)
Expand Down
Loading
Loading