You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何让Step Functions仅在Glue作业抛出RetryableException时重试?

实现Glue作业仅针对自定义RetryableException重试的方案

要实现仅在抛出RetryableException时重试Glue作业,同时将NonRetryableException的消息发送至DLQ,可以通过Glue作业结构化错误输出+Step Functions分支判断的方式实现,具体步骤如下:

1. 修改Glue作业代码,输出结构化错误信息

在Glue作业中捕获自定义异常时,将错误类型以结构化格式(如JSON)输出,让Step Functions能够识别错误类型。以Python作业为例:

import sys
import json

# 自定义异常
class RetryableException(Exception):
    pass

class NonRetryableException(Exception):
    pass

try:
    # 核心业务逻辑示例
    # raise NonRetryableException("数据格式错误,无需重试")
    raise RetryableException("临时网络故障,需要重试")
except RetryableException as e:
    # 输出包含错误类型的结构化信息
    error_info = {"error_type": "RetryableException", "message": str(e)}
    print(f"ERROR: {json.dumps(error_info)}")
    sys.exit(1)
except NonRetryableException as e:
    error_info = {"error_type": "NonRetryableException", "message": str(e)}
    print(f"ERROR: {json.dumps(error_info)}")
    sys.exit(1)

作业执行失败时,上述代码会将错误类型写入作业的输出结果中,供Step Functions读取。

2. 配置Step Functions状态机,实现分支逻辑

在Step Functions中,先捕获Glue作业的通用失败错误Glue.FailedJob,再通过Choice状态判断错误类型,分别执行重试或发送DLQ操作:

状态机关键配置示例(JSON)

{
  "Comment": "Glue作业重试与DLQ处理流程",
  "StartAt": "RunGlueJob",
  "States": {
    "RunGlueJob": {
      "Type": "Task",
      "Resource": "arn:aws:states:::glue:startJobRun.sync",
      "Parameters": {
        "JobName": "Your-Glue-Job-Name"
      },
      "Catch": [
        {
          "ErrorEquals": ["Glue.FailedJob"],
          "Next": "CheckErrorType"
        }
      ],
      "End": true
    },
    "CheckErrorType": {
      "Type": "Choice",
      "Choices": [
        {
          "Variable": "$.JobRunDetails.JobRunOutput",
          "StringMatches": "*\"error_type\": \"RetryableException\"*",
          "Next": "RunGlueJob"
        },
        {
          "Variable": "$.JobRunDetails.JobRunOutput",
          "StringMatches": "*\"error_type\": \"NonRetryableException\"*",
          "Next": "SendToDLQ"
        }
      ],
      "Default": "SendToDLQ"
    },
    "SendToDLQ": {
      "Type": "Task",
      "Resource": "arn:aws:states:::sqs:sendMessage",
      "Parameters": {
        "QueueUrl": "Your-DLQ-Queue-URL",
        "MessageBody": "$"
      },
      "End": true
    }
  }
}

逻辑说明

  • RunGlueJob:同步调用Glue作业,捕获作业失败的通用错误Glue.FailedJob
  • CheckErrorType:通过StringMatches匹配作业输出中的错误类型,判断是重试还是发送DLQ
  • SendToDLQ:调用SQS将错误信息发送至死信队列

3. 补充说明

  • 若Glue作业使用Scala编写,可通过println输出结构化错误信息,逻辑与Python一致
  • 为避免误匹配,建议错误类型的JSON格式保持固定,比如始终用ERROR: {"error_type": "...", ...}的前缀
  • 可在RunGlueJob状态的Retry块中设置重试次数、间隔等参数,避免无限重试

内容的提问来源于stack exchange,提问作者ghostrider

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 04:18:16