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
2 changes: 2 additions & 0 deletions docs/pandas.md
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,8 @@ print(cursor.fetchall())
`to_sql` writes the data as Parquet with pyarrow, so it requires `pip install PyAthena[Pandas,Arrow]`.
Conversion to Parquet and upload to S3 use [ThreadPoolExecutor](https://docs.python.org/3/library/concurrent.futures.html#threadpoolexecutor) by default.
It is also possible to use [ProcessPoolExecutor](https://docs.python.org/3/library/concurrent.futures.html#processpoolexecutor).
The S3 requests use the connection's `s3_config` (see "S3 client" in [Usage](usage.md)) and credentials.
The upload workers resolve the credentials themselves, except for a connection given a `session` and no explicit keys, whose session's credentials they receive once, when the uploads start.

```python
import pandas as pd
Expand Down
3 changes: 3 additions & 0 deletions pyathena/connection.py
Original file line number Diff line number Diff line change
Expand Up @@ -295,6 +295,9 @@ def __init__(
if self.s3_staging_dir and not self.s3_staging_dir.endswith("/"):
self.s3_staging_dir = f"{self.s3_staging_dir}/"

# Whether the session was given rather than built from the arguments,
# which then do not reproduce its credentials.
self._session_given = bool(session)

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 of the repair (both perspectives), scope d764a5917a96e37c67b1de89c64ab79254cc7a8e..4b221020a1139834d813876791f89d38d4955059: CLEAN

Round one (behavior): _session_given is bool(session), matching the if session: branch that uses it. The freeze condition requires a given session, no explicit client keys, no botocore session, and non-empty credentials; every other connection keeps master's worker arguments, so refresh behavior is unchanged there. test_to_sql_session_credentials_and_s3_config (thread and spawn process pools) fails on master's source and passes; test_to_sql_workers_resolve_credentials passes (it would fail on d764a59, which froze for every connection). tests/pyathena/pandas/test_util.py: 15 passed; just lint passed.

Round two (claims): the docstring and docs/pandas.md now say the workers resolve credentials themselves except for a given session; the PR body's behavior changes limit the once-at-start resolution to connect(session=...). Connection._session_given is private and only read by to_sql(). just docs lint passed.

if session:
self._session = session
else:
Expand Down
28 changes: 26 additions & 2 deletions pyathena/pandas/util.py
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,11 @@ def to_sql(
as Parquet files to S3 and executing the appropriate DDL statements.
Supports partitioning, compression, and parallel uploads.

The S3 requests use the connection's ``s3_config`` and credentials. The
upload workers resolve the credentials themselves, except for a connection
given a ``session`` and no explicit keys, whose session's credentials they
receive once, when the uploads start.

Args:
df: The DataFrame to write to Athena.
name: Name of the table to create.
Expand Down Expand Up @@ -259,7 +264,7 @@ def to_sql(

bucket_name, key_prefix = parse_output_location(location)
bucket = conn.session.resource(
"s3", region_name=conn.region_name, **conn._s3_client_kwargs
"s3", region_name=conn.region_name, config=conn.s3_config, **conn._s3_client_kwargs
).Bucket(bucket_name)
cursor = conn.cursor()

Expand Down Expand Up @@ -296,10 +301,29 @@ def to_sql(
reset_index(df, index_label)
with executor_class(max_workers=max_workers) as e:
futures: list[concurrent.futures.Future[Any]] = []
# The workers build their own sessions from these arguments, which
# resolve the connection's credentials again unless the connection was
# given its session. Then they get its credentials as of now, unless
# explicit keys or a botocore session, which boto3 would set them on,
# take precedence.
session_kwargs = deepcopy(conn._session_kwargs)
session_kwargs.update({"profile_name": conn.profile_name})
if (
conn._session_given

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, static review.

  • Reviewer: OpenAI Codex CLI 0.160.0, model gpt-6-astra, effort high, sandbox read-only, session 01a10536-d1ce-7181-8df3-a4a8e53134f8. Scope 911492c3..d764a591 in a detached snapshot without .env; no PR text or earlier findings in the prompt; snapshot unchanged afterwards.
  • Covered (reviewer's list): Connection session construction, credential precedence, _session_kwargs/_s3_client_kwargs/config merging; replacement listing/deletion, partitioned and unpartitioned uploads, worker serialization, retries; both executors, supplied-session side effects, refreshable/missing credentials, tests, docstring and docs; boto3/botocore 1.43.102 and multiprocessing/pickle sources.

Findings and dispositions:

  1. (Introduced, P2) Freezing credentials for every connection made the workers' refreshable credentials (assume-role profiles, instance/container roles) expire in a long upload. Verified; redesigned in 4b22102 after the maintainer chose to freeze only for connect(session=...): Connection._session_given records a given session; otherwise the workers resolve credentials themselves as before. test_to_sql_workers_resolve_credentials checks that a connection without a session passes no credentials to the workers.
  2. (Introduced, P2) s3_config=Config(signature_version=UNSIGNED) cannot be pickled for ProcessPoolExecutor. Verified; not changed: unsigned requests cannot write a table's data, so this configuration cannot work with to_sql() anyway. Stated in the PR body.
  3. (Introduced, P2) The session's credentials were resolved even when explicit keys take precedence, so a partial-credentials environment could raise PartialCredentialsError. Repaired in 4b22102: no lookup when aws_access_key_id is among the client arguments.
  4. (Introduced, P2) Under the forkserver start method (Linux default since Python 3.14), the process-pool test could start a forkserver that keeps the test's invalid keys for later pools. Repaired in 4b22102: the test uses a ProcessPoolExecutor with the spawn context, whose workers exit with the pool.
  5. (Pre-existing) A botocore_session among the connection arguments is deep-copied for the workers, which fails for one holding refreshable credentials and cannot be pickled for processes. Not changed; the code comment no longer claims the arguments are picklable, and the PR body lists it under "Not changed".

and not conn._s3_client_kwargs.get("aws_access_key_id")

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 review (relayed): FINDINGS, static review.

  • Reviewer: OpenAI Codex CLI 0.160.0, model gpt-6-astra, effort high, sandbox read-only, session 01a10541-28d5-7932-8526-98cd168c702a. Scope d764a591..4b221020 (the redesign) in a detached snapshot without .env; snapshot unchanged afterwards.
  • Covered (reviewer's list): default/profile credentials, explicit keys, assume-role/MFA, supplied boto3/botocore sessions, both executors; cleanup and upload S3 configuration incl. partitions; tests, environment restoration, process isolation, docstrings and docs; boto3/botocore 1.43.102. No test environment or monkeypatch leak found; the spawn context avoids the forkserver problem.

Findings and dispositions:

  1. (Introduced, P2) The key-presence check skipped forwarding for aws_access_key_id=None, which botocore treats as unset; SQLAlchemy's create_connect_args() always passes it (None without credentials in the URL). Verified and repaired in ab57b06: both checks look at values. New test case [ThreadPoolExecutor-connect_kwargs2] (aws_access_key_id=None, aws_secret_access_key=None) fails with InvalidAccessKeyId on PutObject with the 4b22102 check and passes now.
  2. (Introduced in d764a59, P2) botocore_session=None likewise suppressed forwarding. Repaired in ab57b06 by the same value check.
  3. (Pre-existing) A botocore session holding refreshable credentials cannot be deep-copied for the workers. Not changed; listed under "Not changed" in the PR body.
  4. Docs wording omitted the explicit-key precedence. Repaired in ab57b06 (docstring and docs/pandas.md: "a session and no explicit keys").

After the repair: just lint, just docs lint, tests/pyathena/pandas/test_util.py: 16 passed.

and not session_kwargs.get("botocore_session")

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 review of ab57b06 (relayed): CLEAN, static review.

  • Reviewer: OpenAI Codex CLI 0.160.0, model gpt-6-astra, effort high, sandbox read-only, session 01a1054a-3c5e-7ce2-b5eb-c951416abdf9. Scope 4b221020..ab57b063 in a detached snapshot without .env; snapshot unchanged afterwards.
  • Result: "No actionable defects introduced by this commit. The value checks handle SQLAlchemy's None credential arguments, and the new test parameter detects the original failure. The wording change introduces no new inaccuracy."

and (credentials := conn.session.get_credentials())
):
frozen_credentials = credentials.get_frozen_credentials()
session_kwargs.update(
{
"aws_access_key_id": frozen_credentials.access_key,
"aws_secret_access_key": frozen_credentials.secret_key,
"aws_session_token": frozen_credentials.token,
}
)
client_kwargs = deepcopy(conn._s3_client_kwargs)
client_kwargs.update({"region_name": conn.region_name})
client_kwargs.update({"region_name": conn.region_name, "config": conn.s3_config})

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): CLEAN

Base 911492c, head d764a59 (PR body, commit messages, docstring, docs/pandas.md, issue #1067).

  • "Proxies, timeouts, retries, max_pool_connections and the user agent did not apply": master's conn.session.resource("s3", region_name=..., **conn._s3_client_kwargs) and session.resource("s3", **client_kwargs) pass no config, and _client_kwargs never contains config (it is a named Connection parameter).
  • "The uploads used the default credential chain with connect(session=...)": reproduced in to_sql() ignores the connection's botocore config and the credentials of connect(session=...)聽#1067 (worker session_kwargs == {'profile_name': None}), and the live test fails on master with the environment's invalid keys.
  • "Credentials resolved once": get_frozen_credentials() is called once per to_sql() call before the chunks are submitted; before, each chunk's new session resolved them again. Stated as a behavior change, with the role_arn/serial_number precedent (fixed temporary keys in _kwargs).
  • Existing callers: to_sql() and to_parquet() signatures are unchanged; to_parquet() callers that pass their own kwargs are unaffected.
  • AWS operator: the S3 requests now follow the connection's retry config inside botocore in addition to retry_api_call around PutObject, as the result-set requests already do.
  • Docs: the docs/pandas.md sentence links to the existing "S3 client" section; just docs lint passed.

partition_prefixes = []
if partitions:
for keys, group in df.groupby(by=partitions, observed=True):
Expand Down
85 changes: 85 additions & 0 deletions tests/pyathena/pandas/test_util.py
Original file line number Diff line number Diff line change
@@ -1,13 +1,20 @@
import contextlib
import textwrap
import uuid
from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor
from datetime import date, datetime
from decimal import Decimal
from multiprocessing import get_context
from unittest.mock import patch

import numpy as np
import pandas as pd
import pytest
from boto3.session import Session
from botocore.config import Config

from pyathena import OperationalError
from pyathena.pandas import util
from pyathena.pandas.util import (
as_pandas,
generate_ddl,
Expand All @@ -16,6 +23,7 @@
to_sql,
)
from tests import ENV
from tests.pyathena.conftest import connect


def test_get_chunks():
Expand Down Expand Up @@ -478,6 +486,83 @@ def test_to_sql_athena_endpoint_url(cursor):
assert cursor.fetchall() == [(1,)]


class SpawnProcessPoolExecutor(ProcessPoolExecutor):
"""Process pool whose workers exit with it and keep no environment for later pools."""

def __init__(self, max_workers=None):
super().__init__(max_workers, mp_context=get_context("spawn"))


@pytest.mark.parametrize(
("executor_class", "connect_kwargs"),
[
(ThreadPoolExecutor, {}),
(SpawnProcessPoolExecutor, {}),
# As SQLAlchemy passes them without credentials in the URL.
(ThreadPoolExecutor, {"aws_access_key_id": None, "aws_secret_access_key": None}),
],
)
def test_to_sql_session_credentials_and_s3_config(monkeypatch, executor_class, connect_kwargs):
# GH-1067: the upload workers used the default credential chain instead of
# the credentials of connect(session=...), and no S3 request used s3_config.
credentials = Session().get_credentials().get_frozen_credentials()
session = Session(
aws_access_key_id=credentials.access_key,
aws_secret_access_key=credentials.secret_key,
aws_session_token=credentials.token,
region_name=ENV.region_name,
)
# The default credential chain, as the workers would resolve it, fails.
for key in ("AWS_PROFILE", "AWS_DEFAULT_PROFILE", "AWS_SESSION_TOKEN"):
monkeypatch.delenv(key, raising=False)
monkeypatch.setenv("AWS_ACCESS_KEY_ID", "AKIAIOSFODNN7EXAMPLE")
monkeypatch.setenv("AWS_SECRET_ACCESS_KEY", "invalid")
df = pd.DataFrame({"col_int": np.int32([1, 2])})
table_name = f"""to_sql_{str(uuid.uuid4()).replace("-", "")}"""
location = f"{ENV.s3_staging_dir}{ENV.schema}/{table_name}/"
resource = Session.resource
with (
contextlib.closing(
connect(
schema_name=ENV.schema,
session=session,
s3_config=Config(max_pool_connections=37),
**connect_kwargs,
)
) as conn,
patch.object(Session, "resource", autospec=True, side_effect=resource) as resources,
):
to_sql(
df,
table_name,
conn,
location,
schema=ENV.schema,
chunksize=1,
executor_class=executor_class,
max_workers=2,
)
cursor = conn.cursor()
cursor.execute(f"SELECT * FROM {table_name} ORDER BY col_int")
assert cursor.fetchall() == [(1,), (2,)]
# The bucket resource, and with threads also the workers' resources.
expected = 1 if executor_class is SpawnProcessPoolExecutor else 3
assert len(resources.call_args_list) == expected
assert all(c.kwargs["config"].max_pool_connections == 37 for c in resources.call_args_list)


def test_to_sql_workers_resolve_credentials(cursor):
# Without a given session, the workers resolve the credentials themselves,
# so refreshable credentials stay refreshable.
df = pd.DataFrame({"col_int": np.int32([1])})
table_name = f"""to_sql_{str(uuid.uuid4()).replace("-", "")}"""
location = f"{ENV.s3_staging_dir}{ENV.schema}/{table_name}/"
with patch.object(util, "to_parquet", side_effect=util.to_parquet) as to_parquet:
to_sql(df, table_name, cursor._connection, location, schema=ENV.schema)
session_kwargs = to_parquet.call_args.args[4]
assert "aws_secret_access_key" not in session_kwargs


def test_to_sql_with_partitions(cursor):
df = pd.DataFrame(
{
Expand Down
Loading