Repository navigation
Use the connection's s3_config and session credentials in to_sql() #1068
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We鈥檒l occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
d540ea5
d764a59
4b22102
ab57b06
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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. | ||
|
|
@@ -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() | ||
|
|
||
|
|
@@ -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 | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent review (relayed): FINDINGS, static review.
Findings and dispositions:
|
||
| and not conn._s3_client_kwargs.get("aws_access_key_id") | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent follow-up review (relayed): FINDINGS, static review.
Findings and dispositions:
After the repair: |
||
| and not session_kwargs.get("botocore_session") | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Independent follow-up review of ab57b06 (relayed): CLEAN, static review.
|
||
| 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}) | ||
|
Member
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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,
|
||
| partition_prefixes = [] | ||
| if partitions: | ||
| for keys, group in df.groupby(by=partitions, observed=True): | ||
|
|
||
There was a problem hiding this comment.
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: CLEANRound one (behavior):
_session_givenisbool(session), matching theif 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_credentialspasses (it would fail on d764a59, which froze for every connection).tests/pyathena/pandas/test_util.py: 15 passed;just lintpassed.Round two (claims): the docstring and
docs/pandas.mdnow say the workers resolve credentials themselves except for a given session; the PR body's behavior changes limit the once-at-start resolution toconnect(session=...).Connection._session_givenis private and only read byto_sql().just docs lintpassed.