From d540ea5a2c91eaa2653c21e0b275b0f8f0898d0f Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 13:38:14 +0900 Subject: [PATCH 1/4] Use the connection's s3_config and session credentials in to_sql() to_sql() built its S3 resources without the connection's botocore config, so proxies, timeouts, retries and the user agent did not apply, and its upload workers built sessions from the connection's keyword arguments only, so the credentials of connect(session=...) were ignored and the uploads used the default credential chain. The resources now get s3_config, and the workers get the frozen credentials of the connection's session, which stay picklable for a ProcessPoolExecutor. Closes #1067 Co-Authored-By: Claude Opus 5.5 --- docs/pandas.md | 1 + pyathena/pandas/util.py | 19 +++++++++-- tests/pyathena/pandas/test_util.py | 55 ++++++++++++++++++++++++++++++ 3 files changed, 73 insertions(+), 2 deletions(-) diff --git a/docs/pandas.md b/docs/pandas.md index bdb0f749..900cb22c 100644 --- a/docs/pandas.md +++ b/docs/pandas.md @@ -112,6 +112,7 @@ 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 the credentials of its session, which the upload workers receive once, when the uploads start. ```python import pandas as pd diff --git a/pyathena/pandas/util.py b/pyathena/pandas/util.py index 579461cc..3e36536f 100644 --- a/pyathena/pandas/util.py +++ b/pyathena/pandas/util.py @@ -215,6 +215,9 @@ 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 the credentials of + its session, which the upload workers receive once, when the uploads start. + Args: df: The DataFrame to write to Athena. name: Name of the table to create. @@ -259,7 +262,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() @@ -296,10 +299,22 @@ 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 picklable arguments, for a + # ProcessPoolExecutor, with the credentials of the connection's session. session_kwargs = deepcopy(conn._session_kwargs) session_kwargs.update({"profile_name": conn.profile_name}) + credentials = conn.session.get_credentials() + if 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}) partition_prefixes = [] if partitions: for keys, group in df.groupby(by=partitions, observed=True): diff --git a/tests/pyathena/pandas/test_util.py b/tests/pyathena/pandas/test_util.py index d74ee430..1051f547 100644 --- a/tests/pyathena/pandas/test_util.py +++ b/tests/pyathena/pandas/test_util.py @@ -1,11 +1,16 @@ +import contextlib import textwrap import uuid +from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor from datetime import date, datetime from decimal import Decimal +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.util import ( @@ -16,6 +21,7 @@ to_sql, ) from tests import ENV +from tests.pyathena.conftest import connect def test_get_chunks(): @@ -478,6 +484,55 @@ def test_to_sql_athena_endpoint_url(cursor): assert cursor.fetchall() == [(1,)] +@pytest.mark.parametrize("executor_class", [ThreadPoolExecutor, ProcessPoolExecutor]) +def test_to_sql_session_credentials_and_s3_config(monkeypatch, executor_class): + # 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), + ) + ) 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 ProcessPoolExecutor 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_with_partitions(cursor): df = pd.DataFrame( { From d764a5917a96e37c67b1de89c64ab79254cc7a8e Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 13:40:00 +0900 Subject: [PATCH 2/4] Leave a given botocore session's credentials to it in to_sql() Co-Authored-By: Claude Opus 5.5 --- pyathena/pandas/util.py | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/pyathena/pandas/util.py b/pyathena/pandas/util.py index 3e36536f..4a641d41 100644 --- a/pyathena/pandas/util.py +++ b/pyathena/pandas/util.py @@ -301,9 +301,13 @@ def to_sql( futures: list[concurrent.futures.Future[Any]] = [] # The workers build their own sessions from picklable arguments, for a # ProcessPoolExecutor, with the credentials of the connection's session. + # A given botocore session already has them, and boto3 would set the + # credentials on it. session_kwargs = deepcopy(conn._session_kwargs) session_kwargs.update({"profile_name": conn.profile_name}) - credentials = conn.session.get_credentials() + credentials = ( + None if "botocore_session" in session_kwargs else conn.session.get_credentials() + ) if credentials: frozen_credentials = credentials.get_frozen_credentials() session_kwargs.update( From 4b221020a1139834d813876791f89d38d4955059 Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 13:51:05 +0900 Subject: [PATCH 3/4] Pass fixed credentials to to_sql()'s workers only for a given session Freezing the credentials for every connection made refreshable credentials (assume-role profiles, instance and container roles) expire in a long upload, and resolved the session's credentials even when explicit keys take precedence. The workers now resolve the credentials themselves as before, unless the connection was given its session; then they get its credentials as of the start of the uploads. Connection records whether its session was given. The process pool test uses the spawn start method, so that no forkserver keeps the invalid keys of its environment for later pools. Co-Authored-By: Claude Opus 5.5 --- docs/pandas.md | 3 ++- pyathena/connection.py | 3 +++ pyathena/pandas/util.py | 25 +++++++++++++++---------- tests/pyathena/pandas/test_util.py | 25 +++++++++++++++++++++++-- 4 files changed, 43 insertions(+), 13 deletions(-) diff --git a/docs/pandas.md b/docs/pandas.md index 900cb22c..5a45d1af 100644 --- a/docs/pandas.md +++ b/docs/pandas.md @@ -112,7 +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 the credentials of its session, which the upload workers receive once, when the uploads start. +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`, whose credentials they receive once, when the uploads start. ```python import pandas as pd diff --git a/pyathena/connection.py b/pyathena/connection.py index 8b1435f9..6d52074b 100644 --- a/pyathena/connection.py +++ b/pyathena/connection.py @@ -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) if session: self._session = session else: diff --git a/pyathena/pandas/util.py b/pyathena/pandas/util.py index 4a641d41..18ebc8f4 100644 --- a/pyathena/pandas/util.py +++ b/pyathena/pandas/util.py @@ -215,8 +215,10 @@ 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 the credentials of - its session, which the upload workers receive once, when the uploads start. + 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``, whose credentials they receive once, when the uploads + start. Args: df: The DataFrame to write to Athena. @@ -299,16 +301,19 @@ 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 picklable arguments, for a - # ProcessPoolExecutor, with the credentials of the connection's session. - # A given botocore session already has them, and boto3 would set the - # credentials on it. + # 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}) - credentials = ( - None if "botocore_session" in session_kwargs else conn.session.get_credentials() - ) - if credentials: + if ( + conn._session_given + and "aws_access_key_id" not in conn._s3_client_kwargs + and "botocore_session" not in session_kwargs + and (credentials := conn.session.get_credentials()) + ): frozen_credentials = credentials.get_frozen_credentials() session_kwargs.update( { diff --git a/tests/pyathena/pandas/test_util.py b/tests/pyathena/pandas/test_util.py index 1051f547..93a5e5c0 100644 --- a/tests/pyathena/pandas/test_util.py +++ b/tests/pyathena/pandas/test_util.py @@ -4,6 +4,7 @@ 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 @@ -13,6 +14,7 @@ from botocore.config import Config from pyathena import OperationalError +from pyathena.pandas import util from pyathena.pandas.util import ( as_pandas, generate_ddl, @@ -484,7 +486,14 @@ def test_to_sql_athena_endpoint_url(cursor): assert cursor.fetchall() == [(1,)] -@pytest.mark.parametrize("executor_class", [ThreadPoolExecutor, ProcessPoolExecutor]) +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", [ThreadPoolExecutor, SpawnProcessPoolExecutor]) def test_to_sql_session_credentials_and_s3_config(monkeypatch, executor_class): # GH-1067: the upload workers used the default credential chain instead of # the credentials of connect(session=...), and no S3 request used s3_config. @@ -528,11 +537,23 @@ def test_to_sql_session_credentials_and_s3_config(monkeypatch, executor_class): 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 ProcessPoolExecutor else 3 + 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( { From ab57b0633b4efba0ae94d85097a6380b4882d2ca Mon Sep 17 00:00:00 2001 From: laughingman7743 Date: Sun, 4 Oct 2026 14:01:09 +0900 Subject: [PATCH 4/4] Check the values of explicit keys and botocore session in to_sql() SQLAlchemy passes aws_access_key_id=None without credentials in the URL, and None does not take precedence over the given session's credentials. Co-Authored-By: Claude Opus 5.5 --- docs/pandas.md | 2 +- pyathena/pandas/util.py | 8 ++++---- tests/pyathena/pandas/test_util.py | 13 +++++++++++-- 3 files changed, 16 insertions(+), 7 deletions(-) diff --git a/docs/pandas.md b/docs/pandas.md index 5a45d1af..68d050ef 100644 --- a/docs/pandas.md +++ b/docs/pandas.md @@ -113,7 +113,7 @@ print(cursor.fetchall()) 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`, whose credentials they receive once, when the uploads start. +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 diff --git a/pyathena/pandas/util.py b/pyathena/pandas/util.py index 18ebc8f4..899a2344 100644 --- a/pyathena/pandas/util.py +++ b/pyathena/pandas/util.py @@ -217,8 +217,8 @@ def to_sql( 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``, whose credentials they receive once, when the uploads - start. + 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. @@ -310,8 +310,8 @@ def to_sql( session_kwargs.update({"profile_name": conn.profile_name}) if ( conn._session_given - and "aws_access_key_id" not in conn._s3_client_kwargs - and "botocore_session" not in session_kwargs + and not conn._s3_client_kwargs.get("aws_access_key_id") + and not session_kwargs.get("botocore_session") and (credentials := conn.session.get_credentials()) ): frozen_credentials = credentials.get_frozen_credentials() diff --git a/tests/pyathena/pandas/test_util.py b/tests/pyathena/pandas/test_util.py index 93a5e5c0..4564f3fc 100644 --- a/tests/pyathena/pandas/test_util.py +++ b/tests/pyathena/pandas/test_util.py @@ -493,8 +493,16 @@ def __init__(self, max_workers=None): super().__init__(max_workers, mp_context=get_context("spawn")) -@pytest.mark.parametrize("executor_class", [ThreadPoolExecutor, SpawnProcessPoolExecutor]) -def test_to_sql_session_credentials_and_s3_config(monkeypatch, executor_class): +@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() @@ -519,6 +527,7 @@ def test_to_sql_session_credentials_and_s3_config(monkeypatch, executor_class): 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,