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
3 changes: 2 additions & 1 deletion docs/filesystem.md
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,8 @@ fsspec.register_implementation("s3", s3fs.S3FileSystem, clobber=True)

## Basic usage

The filesystem can be constructed from a PyAthena connection, or directly with
The filesystem can be constructed from a PyAthena connection, whose S3 client it
then uses (see "S3 client" in [Usage](usage.md)), or directly with
s3fs-compatible credential arguments:

```python
Expand Down
28 changes: 28 additions & 0 deletions docs/usage.md
Original file line number Diff line number Diff line change
Expand Up @@ -750,6 +750,34 @@ The connection builds one Glue client on first use from its session, region, and
Pass `glue_metadata_fallback=False` to `connect()` to turn the fallback off.
The Glue request does not carry the connection's workgroup; turn the fallback off where access depends on the workgroup, such as a workgroup enabled for IAM Identity Center.

## S3 client

A connection builds one S3 client on first use and shares it with the result sets of its cursors, the `S3FileSystem` instances created with `connection=`, and the Spark cursors.
The result sets of `Cursor`, `DictCursor`, and their asynchronous versions do not use it, so these cursors build no S3 client.
`ArrowCursor`, and `PolarsCursor` for Parquet (`unload=True`) and chunked results, read the query results through their libraries' own S3 clients and use the shared client for their other S3 requests.

The S3 client is built from the connection's session, region, and client arguments, with the botocore `config` merged with `s3_config`.
It does not use the connection's `endpoint_url` and `api_version`, which are Athena's, and neither do the S3 requests of `pyathena.pandas.util.to_sql()`.
To send S3 requests to another endpoint, use botocore's [service-specific endpoint settings](https://docs.aws.amazon.com/sdkref/latest/guide/feature-ss-endpoints.html), such as the `AWS_ENDPOINT_URL_S3` environment variable.
Options set in `s3_config` take precedence over those in `config` for the S3 client only:

```python
from botocore.config import Config
from pyathena import connect

conn = connect(
s3_staging_dir="s3://YOUR_S3_BUCKET/path/to/",
region_name="us-west-2",
s3_config=Config(max_pool_connections=50),
)
```

The client keeps up to `max_pool_connections` connections per host for reuse, 10 by default.
Concurrent requests beyond that open more connections, which urllib3 closes after use, logging a "Connection pool is full" warning.

`Connection.close()` closes the network connections of the connection's Athena client, and of its Glue and S3 clients if they were built.
A client used after that opens new connections.

## Environment variables

Support [Boto3 environment variables](https://boto3.amazonaws.com/v1/documentation/api/latest/guide/configuration.html#using-environment-variables).
Expand Down
74 changes: 64 additions & 10 deletions pyathena/connection.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

import logging
import os
import threading
import time
from collections.abc import Callable
from typing import (
Expand Down Expand Up @@ -132,6 +133,7 @@ def __init__(
on_start_query_execution: Callable[[str], None] | None = ...,
on_poll: OnPollCallback | None = ...,
glue_metadata_fallback: bool = ...,
s3_config: Config | None = ...,
**kwargs,
) -> None: ...

Expand Down Expand Up @@ -165,6 +167,7 @@ def __init__(
on_start_query_execution: Callable[[str], None] | None = ...,
on_poll: OnPollCallback | None = ...,
glue_metadata_fallback: bool = ...,
s3_config: Config | None = ...,
**kwargs,
) -> None: ...

Expand Down Expand Up @@ -197,6 +200,7 @@ def __init__(
on_start_query_execution: Callable[[str], None] | None = None,
on_poll: OnPollCallback | None = None,
glue_metadata_fallback: bool = True,
s3_config: Config | None = None,

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): FINDINGS (3, repaired)

Base a4cf604, head 3c666d1 (full pass: PR body, commit messages, docstrings, docs/usage.md, related docs).

Claims checked:

  • "Two S3 clients per pandas/Polars/S3FS query" and the latency table: measured on 2026-10-04 (20 → 1 clients per 10 queries; script times result-set construction to fetchall()).
  • "botocore default 10": botocore.endpoint.MAX_POOL_CONNECTIONS == 10 (botocore 1.43.102).
  • "Requests beyond the pool do not block": botocore's PoolManager sets no block, and urllib3 2.8.0 _put_conn() closes and discards the connection when the queue is full.
  • "botocore re-creates pools after close()": URLLib3Session.close() only clears the pool managers; a HEAD request after close() succeeded locally.
  • "ArrowCursor reads through pyarrow's S3 filesystem": AthenaArrowResultSet._read_csv()/_read_parquet() use self._fs (pyarrow), and only _get_content_length()/_read_data_manifest() use the shared client.
  • "Default cursors build no S3 client": Cursor/AsyncCursor use AthenaResultSet, AioCursor uses AthenaAioResultSet; neither calls the S3 helpers (live test for Cursor).
  • test_executemany skip: unconditional @pytest.mark.skip, so it is skipped on master too.
  • AWS operator: retries unchanged; the S3 client keeps config's retry settings unless s3_config overrides them, and retry_api_call still wraps the result-set requests.

Findings and repairs (543996a):

  1. Existing caller: s3_config was inserted after config, which moved result_reuse_enable and the later parameters by one position for positional callers. It is now the last parameter, as glue_metadata_fallback was added.
  2. Documentation: the pool sentence said urllib3 closes the extra connections "with a warning"; it is a log record from urllib3.connectionpool, not a Python warning. Reworded in docs/usage.md:774 and the PR body.
  3. Documentation: docs/filesystem.md did not say that a filesystem built from a connection uses its S3 client; it now links to the new section (plain page link, since myst_heading_anchors is not enabled).

After the repairs: just lint, just docs lint, tests/pyathena/test_connection.py (48 passed) and tests/pyathena/spark/test_common.py (86 passed).

**kwargs,
) -> None:
"""Initialize a new Athena database connection.
Expand Down Expand Up @@ -243,6 +247,8 @@ def __init__(
answer a throttled table-metadata, table-listing or
database-listing request from the AWS Glue Data Catalog before
retrying it. Defaults to True.
s3_config: Botocore Config options for the S3 client only, such as
``max_pool_connections``. They are merged over ``config``.
**kwargs: Additional arguments passed to boto3 Session and client.

Raises:
Expand Down Expand Up @@ -331,16 +337,16 @@ def __init__(
**self._session_kwargs,
)

if not self.config.user_agent_extra or (
pyathena.user_agent_extra not in self.config.user_agent_extra
):
self.config.user_agent_extra = (
f"{pyathena.user_agent_extra}"
f"{' ' + self.config.user_agent_extra if self.config.user_agent_extra else ''}"
)
self._add_user_agent(self.config)
self.s3_config: Config = self.config.merge(s3_config) if s3_config else self.config
self._add_user_agent(self.s3_config)
self._client = self._session.client(
"athena", region_name=self.region_name, config=self.config, **self._client_kwargs
)
# Built on first use, once per connection even when several threads
# need it at the same time.
self._s3_client_lock = threading.Lock()
self._s3_client: BaseClient | None = None
self._converter = converter
self._formatter = formatter if formatter else DefaultParameterFormatter()
self._retry_config = retry_config if retry_config else RetryConfig()
Expand All @@ -356,6 +362,19 @@ def __init__(
self._session, self.region_name, self.config, self._client_kwargs
)

@staticmethod
def _add_user_agent(config: Config) -> None:
"""Add PyAthena's user agent to a botocore config unless it has it.

Args:
config: The config to update in place.
"""
if not config.user_agent_extra or pyathena.user_agent_extra not in config.user_agent_extra:
config.user_agent_extra = (
f"{pyathena.user_agent_extra}"
f"{' ' + config.user_agent_extra if config.user_agent_extra else ''}"
)

def _assume_role(
self,
profile_name: str | None,
Expand Down Expand Up @@ -477,6 +496,18 @@ def _client_kwargs(self) -> dict[str, Any]:
"""
return {k: v for k, v in self._kwargs.items() if k in self._CLIENT_PASSING_ARGS}

@property
def _s3_client_kwargs(self) -> dict[str, Any]:
"""Get client keyword arguments for S3 client creation.

Returns:
The client keyword arguments without Athena's ``endpoint_url``
and ``api_version``.
"""
return {
k: v for k, v in self._client_kwargs.items() if k not in ("endpoint_url", "api_version")
}

@property
def session(self) -> Session:
"""Get the boto3 session used for AWS API calls.
Expand All @@ -495,6 +526,24 @@ def client(self) -> BaseClient:
"""
return self._client

@property
def s3_client(self) -> BaseClient:

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

Base a4cf604, head 3d6f7b8 (full diff, 10 files).

Covered:

  • Behavior and failure paths: s3_client builds once under the lock (unit test with 8 threads). Every caller (AthenaResultSet._get_content_length()/_read_data_manifest(), S3FileSystem.__init__, SparkBaseCursor.__init__) only reads the client. Nothing in pyathena/ registers event handlers on it or closes it (grep meta.events, _client.close), so sharing cannot leak state between result sets or filesystems. The result sets of Cursor/DictCursor/AsyncCursor/AioCursor (AthenaResultSet, AthenaAioResultSet) never call the two S3 helpers, so they build no client (live test).
  • After a result set is closed, self.connection is None, so the lazy connection.s3_client access would fail. Every caller of the two helpers runs during result-set construction, or reads through self._fs/storage_options built with the same connection, which already fail after close. No new failure path.
  • close(): botocore re-creates its pools when a closed client is used again (checked locally), so cursors used after close() keep working. The Glue and S3 clients are not built just to close them (unit test).
  • s3_config: Config.merge() returns a new object, so the user's s3_config is not mutated; the user agent is added to the merged config (unit test).
  • Data/resource boundaries: no change to conversions or streams. Framework contracts: fsspec behavior is unchanged except for the client's origin; AioS3FileSystem passes connection to its internal S3FileSystem and shares the client too.
  • Simplicity and tests: the two live regression tests assert client-build counts and fail on the original source.

Finding: the Spark comment pyathena/spark/common.py:108 said the client was "Created" before the session; with a shared client it may already exist. Repaired in 3c666d1 (comment only).

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

  • Reviewer: OpenAI Codex CLI 0.160.0, model gpt-6-astra, reasoning effort high, sandbox read-only, session 01a104a6-8b99-78a1-a0b0-1bdcc7363c2a. A different model from the author (Claude).
  • Scope: git diff a4cf60427f3027db084a51883a59cfea871475f4..543996ac2fb1020706b772c1c9649c6ebce5c25f in a detached snapshot worktree at 543996a, with no .env. The prompt had no PR number, description, commit messages or earlier findings. Constraints: no edits, builds, tests, installs, network or GitHub access.
  • Covered (reviewer's list): lazy S3 construction, locking, construction failures, and all previous client-construction sites; client ownership, event handlers, configuration mutation and close calls across result sets, filesystems and Spark; Connection.close() through sync, asyncio, SQLAlchemy and Spark callers, including later client use; s3_config precedence, user agent, and positional/keyword signature compatibility; the standard, pandas, Polars, Arrow and S3FS result paths, AioS3FileSystem, error translation and retries; tests, docs and docstrings against the dependency sources pinned in uv.lock.
  • Result: "No actionable regressions found in the specified diff. The added tests meaningfully distinguish shared/lazy construction from the previous implementation. Close tests verify delegation and avoid constructing unused clients; they do not establish runtime behavior after close."
  • On that limit: use after close() relies on botocore re-creating its pools. That was checked with a local botocore script (HEAD request after client.close()), not with a PyAthena test.
  • The snapshot and the PR worktree were unchanged after the review (HEAD 543996a, clean status).

"""The S3 client shared by the connection's result sets, filesystems and Spark cursors.

It is built on first use from the connection's session, region,
``s3_config`` and client arguments, except Athena's ``endpoint_url``
and ``api_version``.
"""
with self._s3_client_lock:
if self._s3_client is None:
self._s3_client = self._session.client(
"s3",
region_name=self.region_name,
config=self.s3_config,
**self._s3_client_kwargs,
)
return self._s3_client

@property
def retry_config(self) -> RetryConfig:
"""Get the retry configuration for AWS API calls.
Expand Down Expand Up @@ -587,14 +636,19 @@ def cursor(
def close(self) -> None:
"""Close the connection.

Closes the database connection. This method is provided for DB API 2.0
compatibility. Since Athena connections are stateless, this method
currently does not perform any actual cleanup operations.
Closes the network connections of the Athena client and of the Glue and
S3 clients if they were built. A client used after this opens new
network connections.

Note:
This method is called automatically when using the connection
as a context manager (with statement).
"""
self._client.close()
self._glue.close()
with self._s3_client_lock:
if self._s3_client is not None:
self._s3_client.close()

def commit(self) -> None:
"""Commit any pending transaction.
Expand Down
14 changes: 5 additions & 9 deletions pyathena/filesystem/s3.py
Original file line number Diff line number Diff line change
Expand Up @@ -234,9 +234,10 @@ def __init__(
"""Create a filesystem for Amazon S3.

Args:
connection: A PyAthena connection whose session, region, config and
retry policy the S3 client uses. Without one, the client is built
from s3fs-compatible arguments in ``kwargs``.
connection: A PyAthena connection whose S3 client
(``Connection.s3_client``) and retry policy the filesystem uses.
Without one, the client is built from s3fs-compatible arguments
in ``kwargs``.
default_block_size: The block size for reads and writes; defaults to
``DEFAULT_BLOCK_SIZE``.
default_cache_type: The fsspec cache type for reads; defaults to
Expand All @@ -259,12 +260,7 @@ def __init__(
"""
super().__init__(*args, **kwargs)
if connection:
client = connection.session.client(
"s3",
region_name=connection.region_name,
config=connection.config,
**connection._client_kwargs,
)
client = connection.s3_client

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.

Rebase record: rebased from base a4cf604 / head 7c394d1 onto 8bbf8cb (master after #1044, #1055, #1056, #1057, #1060); new head bc0d251.

  • git range-diff: every commit is unchanged except the first, where the conflict with Add S3Core with typed listing entries and list/HEAD operations #1060 (S3Core) was resolved here: S3FileSystem.__init__ now passes connection.s3_client to S3Core instead of building a client from _client_kwargs. S3FileSystem._client is now core.client, so test_s3_filesystem_uses_connection_s3_client still checks the shared client.
  • Upstream contract check: no new S3 client creation from a connection was added upstream (grep '"s3"', _client_kwargs); the s3fs-compatible path (_get_client_compatible_with_s3fs) is unchanged; docs/filesystem.md merged without conflict.
  • After the rebase: just lint, just docs lint, and pytest -n 2 of tests/pyathena/test_connection.py, spark/test_common.py, sqlalchemy/test_base.py::TestAthenaDialect, the new live tests (pandas share/endpoint, default cursor, to_sql), filesystem/test_s3_core.py and filesystem/test_init.py: 193 + 31 passed. Full coverage is left to the AWS CI on this head.

retry_config = connection.retry_config
else:
client = self._get_client_compatible_with_s3fs(**kwargs)
Expand Down
6 changes: 6 additions & 0 deletions pyathena/glue.py
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,12 @@ def client(self) -> BaseClient:
)
return self._client

def close(self) -> None:
"""Close the network connections of the Glue client if it was built."""
with self._lock:
if self._client is not None:
self._client.close()

@staticmethod
def _catalog_request_kwargs(catalog_name: str | None) -> dict[str, str] | None:
# AwsDataCatalog is the caller's default Glue catalog. An S3 Tables
Expand Down
4 changes: 2 additions & 2 deletions pyathena/pandas/util.py
Original file line number Diff line number Diff line change
Expand Up @@ -259,7 +259,7 @@ def to_sql(

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

Expand Down Expand Up @@ -298,7 +298,7 @@ def to_sql(
futures: list[concurrent.futures.Future[Any]] = []
session_kwargs = deepcopy(conn._session_kwargs)
session_kwargs.update({"profile_name": conn.profile_name})
client_kwargs = deepcopy(conn._client_kwargs)
client_kwargs = deepcopy(conn._s3_client_kwargs)

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 one (behavior), follow-up scope: CLEAN

Scope: 041913b..HEAD (3f1f48c test isolation, cacf8d8 Polars docs, 8ec954b to_sql(), 38f36ba SQLAlchemy use_ssl; the last two added at the maintainer's request after the independent review).

  • to_sql(): the bucket resource (:262, DROP + object deletes for if_exists="replace") and the upload workers (:301, passed to to_parquet() and pickled for a ProcessPoolExecutor) get Connection._s3_client_kwargs, a new dict per access with the same keys as _client_kwargs minus endpoint_url/api_version; region_name is still added explicitly. Credentials and verify still apply.
  • use_ssl (pyathena/sqlalchemy/base.py:309): parsed with strtobool like kill_on_interrupt; an invalid value raises ValueError from create_connect_args() as those options do. All dialects (REST, pandas, Arrow, Polars, S3FS and their aio versions) go through _create_connect_args(). A use_ssl given through connect_args is not in the URL query and is untouched.
  • Behavior consequence verified: use_ssl=false in a URL now reaches botocore as False and switches AWS endpoints to http://; Athena's public endpoint did not accept a plain-HTTP connection within 10 s (curl). Stated in the PR body as a behavior change.
  • Tests: test_to_sql_athena_endpoint_url (live) and test_conn_str_use_ssl (offline, REST + aio, false/true) fail on the source without the respective change; the S3 endpoint unit tests pass with hostile AWS_S3_US_EAST_1_REGIONAL_ENDPOINT/AWS_IGNORE_CONFIGURED_ENDPOINT_URLS values.

client_kwargs.update({"region_name": conn.region_name})
partition_prefixes = []
if partitions:
Expand Down
10 changes: 2 additions & 8 deletions pyathena/result_set.py
Original file line number Diff line number Diff line change
Expand Up @@ -104,12 +104,6 @@ def __init__(
self._hints_by_index[k] = v
else:
self._hints_by_name[k.lower()] = v
self._client = connection.session.client(
"s3",
region_name=connection.region_name,
config=connection.config,
**connection._client_kwargs,
)

self._metadata: tuple[dict[str, Any], ...] | None = None
self._column_types: tuple[str, ...] | None = None
Expand Down Expand Up @@ -773,7 +767,7 @@ def _get_content_length(self) -> int:
bucket, key = parse_output_location(self.output_location)
try:
response = retry_api_call(
self._client.head_object,
self.connection.s3_client.head_object,
config=self._retry_config,
logger=_logger,
Bucket=bucket,
Expand All @@ -791,7 +785,7 @@ def _read_data_manifest(self) -> list[str]:
bucket, key = parse_output_location(self.data_manifest_location)
try:
response = retry_api_call(
self._client.get_object,
self.connection.s3_client.get_object,
config=self._retry_config,
logger=_logger,
Bucket=bucket,
Expand Down
11 changes: 3 additions & 8 deletions pyathena/spark/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -105,14 +105,9 @@ def __init__(
self._calculation_id: str | None = None
self._calculation_execution: AthenaCalculationExecution | None = None

# Created before the session so that a local failure cannot leave
# a newly started session behind.
self._client = self.connection.session.client(
"s3",
region_name=self.connection.region_name,
config=self.connection.config,
**self.connection._client_kwargs,
)
# Taken before the session so that a failure to build the client
# cannot leave a newly started session behind.
self._client = self.connection.s3_client

if session_id:
if self._exists_session(session_id):
Expand Down
2 changes: 2 additions & 0 deletions pyathena/sqlalchemy/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -306,6 +306,8 @@ def _create_connect_args(self, url: URL) -> dict[str, Any]:
with contextlib.suppress(ValueError):
verify = bool(strtobool(verify))
opts.update({"verify": verify})
if "use_ssl" in opts:

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), follow-up scope: CLEAN

Scope: 041913b..HEAD.

  • "botocore treats the string \"false\" as true": Session.client("athena", use_ssl="false").meta.endpoint_url is https://athena.us-west-2.amazonaws.com, and with False it is http://... (botocore 1.43.102, local check).
  • "An Athena VPC endpoint received to_sql()'s PutObject requests": test_to_sql_athena_endpoint_url, with Athena's regional endpoint as endpoint_url, fails without 8ec954b with ClientError ... when calling the PutObject operation.
  • Docs: docs/usage.md now names to_sql() next to the shared client, and the Arrow/Polars sentence matches AthenaPolarsResultSet._parquet_storage_options and the scan_csv/scan_parquet chunk paths, which use Polars' object_store. There is no URL option list in docs/sqlalchemy.md to update.
  • Existing callers: to_sql() signature unchanged; URL users of use_ssl=false are the affected group, documented in the PR body.

opts.update({"use_ssl": bool(strtobool(opts["use_ssl"]))})
if "duration_seconds" in opts:
opts.update({"duration_seconds": int(opts["duration_seconds"])})
if "poll_interval" in opts:
Expand Down
24 changes: 24 additions & 0 deletions tests/pyathena/pandas/test_cursor.py
Original file line number Diff line number Diff line change
Expand Up @@ -268,6 +268,30 @@ def test_result_set_file_system(self, pandas_cursor, chunksize):
if not pandas_cursor.result_set.is_unload:
assert pandas_cursor.result_set._csv_stream.closed

@pytest.mark.parametrize(
"pandas_cursor",
[{"endpoint_url": f"https://athena.{ENV.region_name}.amazonaws.com"}],
indirect=True,
)
def test_athena_endpoint_url(self, pandas_cursor):
# GH-576: the S3 requests were sent to Athena's endpoint_url.
pandas_cursor.execute("SELECT * FROM one_row")
assert pandas_cursor.fetchall() == [(1,)]

def test_result_sets_share_s3_client(self, pandas_cursor):
# GH-1011: each result set and its filesystem built their own S3 client,
# so every query opened new connections to S3.
conn = pandas_cursor.connection
session_client = conn.session.client
file_systems = []
with patch.object(conn.session, "client", side_effect=session_client) as client:
for _ in range(2):
pandas_cursor.execute("SELECT * FROM one_row")
assert pandas_cursor.fetchall() == [(1,)]
file_systems.append(pandas_cursor.result_set._fs)
assert [c.args for c in client.call_args_list] == [("s3",)]
assert all(fs._client is conn.s3_client for fs in file_systems)

@pytest.mark.parametrize(
("query", "expected", "binary"),
[
Expand Down
15 changes: 15 additions & 0 deletions tests/pyathena/pandas/test_util.py
Original file line number Diff line number Diff line change
Expand Up @@ -463,6 +463,21 @@ def test_to_sql_with_index(cursor):
]


@pytest.mark.parametrize(
"cursor",
[{"endpoint_url": f"https://athena.{ENV.region_name}.amazonaws.com"}],
indirect=True,
)
def test_to_sql_athena_endpoint_url(cursor):
# GH-576: the uploads were sent to Athena's endpoint_url.
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}/"
to_sql(df, table_name, cursor._connection, location, schema=ENV.schema, if_exists="fail")
cursor.execute(f"SELECT * FROM {table_name}")
assert cursor.fetchall() == [(1,)]


def test_to_sql_with_partitions(cursor):
df = pd.DataFrame(
{
Expand Down
6 changes: 4 additions & 2 deletions tests/pyathena/spark/test_common.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
import logging
import threading
import uuid
from unittest.mock import MagicMock, patch
from unittest.mock import MagicMock, PropertyMock, patch

import pytest
from botocore.exceptions import ClientError
Expand Down Expand Up @@ -480,7 +480,9 @@ def test_init_does_not_terminate_supplied_session(self, cursor_class):
@pytest.mark.parametrize("cursor_class", SPARK_CURSOR_CLASSES)
def test_init_does_not_start_session_when_s3_client_fails(self, cursor_class):
connection = _connection()
connection.session.client.side_effect = ValueError("Invalid S3 client configuration.")
type(connection).s3_client = PropertyMock(
side_effect=ValueError("Invalid S3 client configuration.")
)
with pytest.raises(ValueError, match=r"^Invalid S3 client configuration\.$"):
_init_cursor(cursor_class, connection)

Expand Down
Loading
Loading