This is an automated email from the ASF dual-hosted git repository.
vincbeck 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 9b2cc1352c3 Handle plain Step Functions execution errors (#73881)
9b2cc1352c3 is described below
commit 9b2cc1352c3716cde9b242473c4ebcb9243a59b0
Author: jingi723 <[email protected]>
AuthorDate: Mon Oct 5 23:00:19 2026 +0900
Handle plain Step Functions execution errors (#73881)
---
providers/amazon/docs/operators/step_functions.rst | 11 +++++
.../amazon/aws/operators/step_function.py | 5 ++-
.../amazon/aws/operators/test_step_function.py | 51 ++++++++++++++++++++++
3 files changed, 66 insertions(+), 1 deletion(-)
diff --git a/providers/amazon/docs/operators/step_functions.rst
b/providers/amazon/docs/operators/step_functions.rst
index 4694da5f307..ce7667286d0 100644
--- a/providers/amazon/docs/operators/step_functions.rst
+++ b/providers/amazon/docs/operators/step_functions.rst
@@ -59,6 +59,17 @@ Get an AWS Step Functions execution output
To fetch the output from an AWS Step Function state machine execution you can
use
:class:`~airflow.providers.amazon.aws.operators.step_function.StepFunctionGetExecutionOutputOperator`.
+The operator decodes JSON execution output. When the execution has an error
instead,
+it returns the error string unchanged if it is not valid JSON. JSON-formatted
errors
+continue to be decoded to their corresponding Python values.
+
+Fetching an error does not by itself fail the Airflow task. This operator does
not wait
+for completion or check that the execution succeeded; use
+:class:`~airflow.providers.amazon.aws.sensors.step_function.StepFunctionExecutionSensor`
+to check the execution state. A returned error string is not suitable for
+``multiple_outputs=True`` or downstream dynamic task mapping, which require a
mapping
+or a mappable collection, respectively.
+
.. exampleinclude::
/../../amazon/tests/system/amazon/aws/example_step_functions.py
:language: python
:dedent: 4
diff --git
a/providers/amazon/src/airflow/providers/amazon/aws/operators/step_function.py
b/providers/amazon/src/airflow/providers/amazon/aws/operators/step_function.py
index 4ae43673290..47266a5482a 100644
---
a/providers/amazon/src/airflow/providers/amazon/aws/operators/step_function.py
+++
b/providers/amazon/src/airflow/providers/amazon/aws/operators/step_function.py
@@ -199,7 +199,10 @@ class
StepFunctionGetExecutionOutputOperator(AwsBaseOperator[StepFunctionHook]):
if "output" in execution_status:
response = json.loads(execution_status["output"])
elif "error" in execution_status:
- response = json.loads(execution_status["error"])
+ try:
+ response = json.loads(execution_status["error"])
+ except json.JSONDecodeError:
+ response = execution_status["error"]
self.log.info("Got State Machine Execution output for %s",
self.execution_arn)
diff --git
a/providers/amazon/tests/unit/amazon/aws/operators/test_step_function.py
b/providers/amazon/tests/unit/amazon/aws/operators/test_step_function.py
index 9d58968a702..31d1aecc10e 100644
--- a/providers/amazon/tests/unit/amazon/aws/operators/test_step_function.py
+++ b/providers/amazon/tests/unit/amazon/aws/operators/test_step_function.py
@@ -17,10 +17,16 @@
# under the License.
from __future__ import annotations
+import json
+from contextlib import closing
+from datetime import datetime, timezone
from unittest import mock
+import boto3
import pytest
+from botocore.stub import Stubber
+from airflow.providers.amazon.aws.hooks.step_function import StepFunctionHook
from airflow.providers.amazon.aws.operators.step_function import (
StepFunctionGetExecutionOutputOperator,
StepFunctionStartExecutionOperator,
@@ -112,6 +118,51 @@ class TestStepFunctionGetExecutionOutputOperator:
)
validate_template_fields(operator)
+ @pytest.mark.parametrize(
+ ("error", "expected_output"),
+ [
+ pytest.param("States.TaskFailed", "States.TaskFailed",
id="plain_error"),
+ pytest.param("", "", id="empty_error"),
+ pytest.param("오류", "오류", id="unicode_error"),
+ pytest.param("{message", "{message", id="non_json_error"),
+ pytest.param('{"message": "failed"}', {"message": "failed"},
id="json_object_error"),
+ pytest.param('"failed"', "failed", id="json_string_error"),
+ pytest.param("42", 42, id="json_number_error"),
+ pytest.param("null", None, id="json_null_error"),
+ ],
+ )
+ def test_execute_describe_execution_error(self, error, expected_output):
+ client = boto3.client(
+ "stepfunctions",
+ region_name="us-east-1",
+ aws_access_key_id="testing",
+ aws_secret_access_key="testing",
+ )
+ op = StepFunctionGetExecutionOutputOperator(
+ task_id=self.TASK_ID, execution_arn=EXECUTION_ARN,
aws_conn_id=None, region_name="us-east-1"
+ )
+ with closing(client), Stubber(client) as stubber,
mock.patch.object(op.hook, "conn", client):
+ stubber.add_response(
+ "describe_execution",
+ {
+ "executionArn": EXECUTION_ARN,
+ "stateMachineArn": STATE_MACHINE_ARN,
+ "startDate": datetime(2026, 1, 1, tzinfo=timezone.utc),
+ "status": "FAILED",
+ "error": error,
+ },
+ {"executionArn": EXECUTION_ARN},
+ )
+ assert op.execute({}) == expected_output
+ stubber.assert_no_pending_responses()
+
+ @mock.patch.object(StepFunctionGetExecutionOutputOperator, "hook",
spec=StepFunctionHook)
+ def test_execute_rejects_invalid_json_output(self, mocked_hook):
+ mocked_hook.describe_execution.return_value = {"output": "not-json"}
+ op = StepFunctionGetExecutionOutputOperator(task_id=self.TASK_ID,
execution_arn=EXECUTION_ARN)
+ with pytest.raises(json.JSONDecodeError):
+ op.execute({})
+
class TestStepFunctionStartExecutionOperator:
TASK_ID = "step_function_start_execution_task"