Skip to content

Commit

Permalink
feat: add dtype argument to delete_from_iceberg (#3099)
Browse files Browse the repository at this point in the history
  • Loading branch information
jaidisido authored Feb 19, 2025
1 parent a2b30f7 commit 7cd4d57
Show file tree
Hide file tree
Showing 2 changed files with 11 additions and 0 deletions.
6 changes: 6 additions & 0 deletions awswrangler/athena/_write_iceberg.py
Original file line number Diff line number Diff line change
Expand Up @@ -667,6 +667,7 @@ def delete_from_iceberg_table(
workgroup: str = "primary",
encryption: str | None = None,
kms_key: str | None = None,
dtype: dict[str, str] | None = None,
boto3_session: boto3.Session | None = None,
s3_additional_kwargs: dict[str, Any] | None = None,
catalog_id: str | None = None,
Expand Down Expand Up @@ -702,6 +703,10 @@ def delete_from_iceberg_table(
Valid values: [``None``, ``"SSE_S3"``, ``"SSE_KMS"``]. Notice: ``"CSE_KMS"`` is not supported.
kms_key
For SSE-KMS, this is the KMS key ARN or ID.
dtype
Dictionary of columns names and Athena/Glue types to be casted.
Useful when you have columns with undetermined or mixed data types.
(e.g. {'col name': 'bigint', 'col2 name': 'int'})
boto3_session
The default boto3 session will be used if **boto3_session** receive ``None``.
s3_additional_kwargs
Expand Down Expand Up @@ -763,6 +768,7 @@ def delete_from_iceberg_table(
boto3_session=boto3_session,
s3_additional_kwargs=s3_additional_kwargs,
catalog_id=catalog_id,
dtype=dtype,
index=False,
)

Expand Down
5 changes: 5 additions & 0 deletions tests/unit/test_athena_iceberg.py
Original file line number Diff line number Diff line change
Expand Up @@ -870,6 +870,7 @@ def test_athena_delete_from_iceberg_table(
"id": [1, 2, 3],
"name": ["a", "b", "c"],
"ts": [ts("2020-01-01 00:00:00.0"), ts("2020-01-02 00:00:01.0"), ts("2020-01-03 00:00:00.0")],
"empty": [pd.NA, pd.NA, pd.NA],
}
)
df["id"] = df["id"].astype("Int64") # Cast as nullable int64 type
Expand All @@ -883,6 +884,7 @@ def test_athena_delete_from_iceberg_table(
temp_path=path2,
partition_cols=partition_cols,
keep_files=False,
dtype={"empty": "string"},
)

wr.athena.delete_from_iceberg_table(
Expand All @@ -892,6 +894,7 @@ def test_athena_delete_from_iceberg_table(
temp_path=path2,
merge_cols=["id"],
keep_files=False,
dtype={"empty": "string"},
)

df_actual = wr.athena.read_sql_query(
Expand All @@ -906,10 +909,12 @@ def test_athena_delete_from_iceberg_table(
"id": [3],
"name": ["c"],
"ts": [ts("2020-01-03 00:00:00.0")],
"empty": [pd.NA],
}
)
df_expected["id"] = df_expected["id"].astype("Int64") # Cast as nullable int64 type
df_expected["name"] = df_expected["name"].astype("string")
df_expected["empty"] = df_expected["empty"].astype("string")

assert_pandas_equals(df_expected, df_actual)

Expand Down

0 comments on commit 7cd4d57

Please sign in to comment.