This is an automated email from the ASF dual-hosted git repository.
henry3260 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 7b4564201fa Fail Glue job tasks stopped mid-run in deferrable verbose
mode (#71495)
7b4564201fa is described below
commit 7b4564201fa89813261bc5855a1f1cd9c86e93c8
Author: rjgoyln <[email protected]>
AuthorDate: Fri Aug 14 02:43:41 2026 +0800
Fail Glue job tasks stopped mid-run in deferrable verbose mode (#71495)
---
.../amazon/src/airflow/providers/amazon/aws/triggers/glue.py | 5 +++--
providers/amazon/tests/unit/amazon/aws/triggers/test_glue.py | 9 +++++----
2 files changed, 8 insertions(+), 6 deletions(-)
diff --git a/providers/amazon/src/airflow/providers/amazon/aws/triggers/glue.py
b/providers/amazon/src/airflow/providers/amazon/aws/triggers/glue.py
index e499763aafa..478ce00371a 100644
--- a/providers/amazon/src/airflow/providers/amazon/aws/triggers/glue.py
+++ b/providers/amazon/src/airflow/providers/amazon/aws/triggers/glue.py
@@ -144,7 +144,8 @@ class GlueJobCompleteTrigger(AwsBaseWaiterTrigger):
)
return
- if job_run_state in ("FAILED", "TIMEOUT"):
+ # STOPPED means the run was cancelled before it produced its
output, not that it succeeded.
+ if job_run_state in ("FAILED", "TIMEOUT", "STOPPED"):
yield TriggerEvent(
{
"status": "error",
@@ -154,7 +155,7 @@ class GlueJobCompleteTrigger(AwsBaseWaiterTrigger):
}
)
return
- if job_run_state in ("SUCCEEDED", "STOPPED"):
+ if job_run_state == "SUCCEEDED":
self.log.info(
"Exiting Job %s Run %s State: %s",
self.job_name,
diff --git a/providers/amazon/tests/unit/amazon/aws/triggers/test_glue.py
b/providers/amazon/tests/unit/amazon/aws/triggers/test_glue.py
index 61a1f8385e6..7281fa5ce94 100644
--- a/providers/amazon/tests/unit/amazon/aws/triggers/test_glue.py
+++ b/providers/amazon/tests/unit/amazon/aws/triggers/test_glue.py
@@ -171,13 +171,14 @@ class TestGlueJobTrigger:
assert logs_client.get_log_events.call_count >= 2
@pytest.mark.asyncio
+ @pytest.mark.parametrize("job_run_state", ["FAILED", "TIMEOUT", "STOPPED"])
@mock.patch.object(AwsLogsHook, "get_async_conn")
@mock.patch.object(GlueJobHook, "get_async_conn")
- async def test_verbose_run_job_failed(self, mock_glue_conn,
mock_logs_conn):
- """When verbose=True and the job fails, the trigger yields an error
event."""
+ async def test_verbose_run_failure_states(self, mock_glue_conn,
mock_logs_conn, job_run_state):
+ """When verbose=True and the run ends in a failure state, the trigger
yields an error event."""
glue_client = AsyncMock()
glue_client.get_job_run = AsyncMock(
- return_value={"JobRun": {"JobRunState": "FAILED", "LogGroupName":
"/aws-glue/python-jobs"}}
+ return_value={"JobRun": {"JobRunState": job_run_state,
"LogGroupName": "/aws-glue/python-jobs"}}
)
mock_glue_conn.return_value.__aenter__ =
AsyncMock(return_value=glue_client)
mock_glue_conn.return_value.__aexit__ = AsyncMock(return_value=False)
@@ -198,7 +199,7 @@ class TestGlueJobTrigger:
generator = trigger.run()
event = await generator.asend(None)
assert event.payload["status"] == "error"
- assert "FAILED" in event.payload["message"]
+ assert job_run_state in event.payload["message"]
assert event.payload["run_id"] == "jr_123"
@pytest.mark.asyncio