This is an automated email from the ASF dual-hosted git repository.
shahar1 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new da5af33fd02 Create the Athena Spark work group and results bucket in
the system test (#73929)
da5af33fd02 is described below
commit da5af33fd024530234b20afd5b116a3e965d45bc
Author: Sean Ghaeli <[email protected]>
AuthorDate: Sat Oct 3 01:13:59 2026 -0700
Create the Athena Spark work group and results bucket in the system test
(#73929)
The work group was preconfigured infrastructure, which meant the test could
not
run against a fresh account and two runs shared one work group. Athena
creates a
work group synchronously and it has no CREATING state, so the test can own
it
without a waiter.
The results bucket is created and force-deleted by the test too, like the
other
AWS system tests, and is used as the work group's output location. Athena
rejects a PySpark work group without an ExecutionRole, so the execution role
stays preconfigured and the context now supplies only its ARN.
---
.../system/amazon/aws/example_athena_spark.py | 52 +++++++++++++++++++---
1 file changed, 45 insertions(+), 7 deletions(-)
diff --git a/providers/amazon/tests/system/amazon/aws/example_athena_spark.py
b/providers/amazon/tests/system/amazon/aws/example_athena_spark.py
index d7ccb624127..a1e08050d0a 100644
--- a/providers/amazon/tests/system/amazon/aws/example_athena_spark.py
+++ b/providers/amazon/tests/system/amazon/aws/example_athena_spark.py
@@ -23,6 +23,7 @@ import boto3
from airflow.providers.amazon.aws.hooks.athena import AthenaHook
from airflow.providers.amazon.aws.operators.athena_spark import
AthenaSparkOperator
+from airflow.providers.amazon.aws.operators.s3 import S3CreateBucketOperator,
S3DeleteBucketOperator
from tests_common.test_utils.version_compat import AIRFLOW_V_3_0_PLUS
@@ -40,15 +41,34 @@ except ImportError:
# Compatibility for Airflow < 3.1
from airflow.utils.trigger_rule import TriggerRule # type:
ignore[no-redef,attr-defined]
-from system.amazon.aws.utils import SystemTestContextBuilder
+from system.amazon.aws.utils import ENV_ID_KEY, SystemTestContextBuilder
DAG_ID = "example_athena_spark"
-# The Spark workgroup is preconfigured test infrastructure; this DAG creates
only the session.
-# Test runners can override the default by exporting ATHENA_SPARK_WORK_GROUP.
-ATHENA_SPARK_WORK_GROUP_KEY = "ATHENA_SPARK_WORK_GROUP"
+# Athena rejects a PySpark work group without an execution role, so the role
is preconfigured
+# test infrastructure. The results bucket and the work group are created here.
+EXECUTION_ROLE_ARN_KEY = "EXECUTION_ROLE_ARN"
-sys_test_context_task =
SystemTestContextBuilder().add_variable(ATHENA_SPARK_WORK_GROUP_KEY).build()
+sys_test_context_task =
SystemTestContextBuilder().add_variable(EXECUTION_ROLE_ARN_KEY).build()
+
+
+@task
+def create_work_group(work_group: str, execution_role_arn: str, bucket_name:
str) -> None:
+ client = boto3.client("athena")
+ client.create_work_group(
+ Name=work_group,
+ Configuration={
+ "ExecutionRole": execution_role_arn,
+ "ResultConfiguration": {"OutputLocation": f"s3://{bucket_name}/"},
+ "EngineVersion": {"SelectedEngineVersion": "PySpark engine version
3"},
+ },
+ )
+
+
+@task(trigger_rule=TriggerRule.ALL_DONE)
+def delete_work_group(work_group: str) -> None:
+ client = boto3.client("athena")
+ client.delete_work_group(WorkGroup=work_group, RecursiveDeleteOption=True)
@task
@@ -83,9 +103,16 @@ with DAG(
catchup=False,
) as dag:
test_context = sys_test_context_task()
- athena_spark_work_group = test_context[ATHENA_SPARK_WORK_GROUP_KEY]
+ env_id = test_context[ENV_ID_KEY]
+
+ work_group = f"{env_id}-athena-spark"
+ bucket_name = f"{env_id}-athena-spark-bucket"
+
+ create_bucket = S3CreateBucketOperator(task_id="create_bucket",
bucket_name=bucket_name)
- session_id = start_athena_spark_session(athena_spark_work_group)
+ setup_work_group = create_work_group(work_group,
test_context[EXECUTION_ROLE_ARN_KEY], bucket_name)
+
+ session_id = start_athena_spark_session(work_group)
idle_session_id = wait_for_athena_spark_session(session_id)
# [START howto_operator_athena_spark]
@@ -100,15 +127,26 @@ with DAG(
stop_session = stop_athena_spark_session(session_id)
+ delete_bucket = S3DeleteBucketOperator(
+ task_id="delete_bucket",
+ bucket_name=bucket_name,
+ force_delete=True,
+ trigger_rule=TriggerRule.ALL_DONE,
+ )
+
chain(
# TEST SETUP
test_context,
+ create_bucket,
+ setup_work_group,
session_id,
idle_session_id,
# TEST BODY
run_spark_calculation,
# TEST TEARDOWN
stop_session,
+ delete_work_group(work_group),
+ delete_bucket,
)
from tests_common.test_utils.watcher import watcher