【问题标题】:How to capture output of a parallel state machine in AWS Step function如何在 AWS Step 函数中捕获并行状态机的输出
【发布时间】:2018-12-27 00:17:42
【问题描述】:

我需要运行一个运行并行状态机的 AWS Step 函数,比如两个状态机。我的要求是检查并行机的最终执行状态,如果有任何故障,调用SNS服务发送电子邮件。相当标准的东西,但对于我的生活,我无法弄清楚如何捕捉并行步进机的组合误差。此示例并行机运行

  1. “passtask”只是一个简单的 lambda 传递函数, 和
  2. 运行一个具有 5 秒睡眠计时器的失败任务,并且应该在 5 秒后失败。

如果我执行这台机器,这台机器正确显示 passtask 为成功,failtask 为取消,整体并行任务成功(?????),通知失败任务为取消,状态机的整体执行为“失败”也是。

我希望看到 passtask 成功,失败任务失败,整体并行任务失败,通知失败任务成功。

{
  "Comment": "Parallel Example",
  "StartAt": "Parallel Task",
  "TimeoutSeconds": 120,
  "States": {
    "Parallel Task": {
      "Type": "Parallel",
      "Branches": [
       {
         "StartAt": "passtask",
         "States": {
           "passtask": {
             "Type": "Task",
             "Resource":"arn:xxxxxxxxxxxxxxx:function:passfunction",
             "End": true
           }
         }
       },
       {
         "StartAt": "failtask",
         "States": {
           "failtask": {
             "Type": "Task",
             "Resource":"arn: xxxxxxxxxxxxxxx:function:failfunction",
             "End": true
           }
         }
       }
      ],
      "ResultPath": "$.status",
      "Catch": [
        {
          "ErrorEquals": ["States.ALL"],
          "Next": "Notify Failure"
        }
      ],
      "Next": "Notify Success"
    },
    "Notify Failure": {
      "Type": "Pass",
      "InputPath": "$.input.Cause",
      "End": true
    },
    "Notify Success": {
      "Type": "Pass",
      "Result": "This is a fallback from a task success",
      "End": true
    }
  }
}

【问题讨论】:

  • 请使用code block重新格式化您的示例。
  • 嗯,这个 sn-p 正确显示在我这边。我使用了代码块(那些`符号)。空白的“灰色”空间是 ASL 缩进代码的方式
  • 所以想出了一个问题 - 通知失败需要是类型:“失败”而不是通过。这会使它变红。所以仍然存在两个问题 - 失败任务显示为已取消,整体并行状态显示为成功
  • 如果您点击了我在评论中发布的链接,您就会知道 代码块 是通过将每行缩进四个空格和 内联代码 用反引号完成。请使用 code block 重新格式化您的示例。

标签: amazon-web-services


【解决方案1】:

从您的要求“我的要求是检查并行机的最终执行状态,如果有任何故障,调用SNS服务发送电子邮件。”,我理解“failtask”只是为了调试目的和将来它不一定会失败。所以问题是,当 Step Functions 检测到一个分支中的故障时,所有其他分支都被终止并且它们的输出被丢弃,只有故障分支的输出被使用。因此,如果要保留每个分支的输出并检查是否发生故障,则需要处理每个分支中的错误,而不是将整个分支报告为失败。此外,您需要向每个分支添加一个输出字段,说明是否存在故障(如果字段不存在,Choice State 将给出错误)。还要记住,ParralelState 的输出是一个包含每个分支输出的数组,例如这个状态机应该让每个分支完成执行并正确处理错误:

{
    "Comment": "Parallel Example",
    "StartAt": "Parallel Task",
    "TimeoutSeconds": 120,
    "States": {
        "Parallel Task": {
            "Type": "Parallel",
            "Branches": [{
                    "StartAt": "passtask",
                    "States": {
                        "passtask": {
                            "Type": "Task",
                            "Resource": "arn:aws:lambda:us-east-1:XXXXXXXXXXXXXXXXX",
                            "Next": "SuccessBranch1",
                            "Catch": [{
                                "ErrorEquals": ["States.ALL"],
                                "Next": "FailBranch1"
                            }]
                        },
                        "SuccessBranch1": {
                            "Type": "Pass",
                            "Result": {
                                "Error": false
                            },
                            "ResultPath": "$.Status",
                            "End": true
                        },
                        "FailBranch1": {
                            "Type": "Pass",
                            "Result": {
                                "Error": true
                            },
                            "ResultPath": "$.Status",
                            "End": true
                        }
                    }

                },
                {
                    "StartAt": "failtask",
                    "States": {
                        "failtask": {
                            "Type": "Task",
                            "Resource": "arn:aws:lambda:us-east-1:XXXXXXXXXXXXXXXXX",
                            "Next": "SuccessBranch2",
                            "Catch": [{
                                "ErrorEquals": ["States.ALL"],
                                "Next": "FailBranch2"
                            }]
                        },
                        "SuccessBranch2": {
                            "Type": "Pass",
                            "Result": {
                                "Error": false
                            },
                            "ResultPath": "$.Status",
                            "End": true
                        },
                        "FailBranch2": {
                            "Type": "Pass",
                            "Result": {
                                "Error": true
                            },
                            "ResultPath": "$.Status",
                            "End": true
                        }
                    }
                }
            ],
            "ResultPath": "$.ParralelOutput",
            "Catch": [{
                "Comment": "This catch should never catch any errors, as the error handling is done in the individual Branches",
                "ErrorEquals": ["States.ALL"],
                "ResultPath": "$.ParralelOutput",
                "Next": "ChoiceStateX"
            }],
            "Next": "ChoiceStateX"
        },

        "ChoiceStateX": {
            "Type": "Choice",

            "Choices": [{
                "Or": [{
                        "Variable": "$.ParralelOutput[0].Status.Error",
                        "BooleanEquals": true
                    },
                    {
                        "Variable": "$.ParralelOutput[1].Status.Error",
                        "BooleanEquals": true
                    }
                ],
                "Next": "Notify Failure"
            }],
            "Default": "Notify Success"
        },

        "Notify Failure": {
            "Type": "Pass",
            "End": true
        },
        "Notify Success": {
            "Type": "Pass",
            "Result": "This is a fallback from a task success",
            "End": true
        }
    }
}

对于 Nisman 在 cmets 中提出的上述更一般的情况(尽管更复杂)。我们可以添加一个带有一些 JSONPath 技巧的通过状态,而不是硬编码选择状态来检查每个分支,以检查当前仅使用选择状态无法实现的条件。

在这个 Pass State 中,我们使用参数来重组我们的数据,这样当我们使用 OutputPath 对这个数据应用 JSONPath 过滤器表达式时,我们会得到一个 2(如果没有分支失败)或 3(如果某些分支失败)元素,其中第一个元素始终包含原始输入数据,第二个/第三个元素至少包含一个同名的键,供选择状态使用。这是状态机 JSON:

{
  "Comment": "Parallel Example",
  "StartAt": "Parallel Task",
  "States": {
    "Parallel Task": {
      "Type": "Parallel",
      "Branches": [
        {
          "StartAt": "passtask",
          "States": {
            "passtask": {
              "Type": "Task",
              "Resource": "<TASK RESOURCE>",
              "End": true,
              "Catch": [
                {
                  "ErrorEquals": [
                    "States.ALL"
                  ],
                  "ResultPath": "$.error-info",
                  "Next": "FailBranch1"
                }
              ]
            },
            "FailBranch1": {
              "Type": "Pass",
              "Parameters": {
                "BranchOutput.$": "$",
                "BranchError": true
              },
              "End": true
            }
          }
        },
        {
          "StartAt": "failtask",
          "States": {
            "failtask": {
              "Type": "Task",
              "Resource": "<TASK RESOURCE>",
              "End": true,
              "Catch": [
                {
                  "ErrorEquals": [
                    "States.ALL"
                  ],
                  "ResultPath": "$.error-info",
                  "Next": "FailBranch2"
                }
              ]
            },
            "FailBranch2": {
              "Type": "Pass",
              "Parameters": {
                "BranchOutput.$": "$",
                "BranchError": true
              },
              "End": true
            }
          }
        }
      ],
      "ResultPath": "$.ParralelOutput",
      "Next": "Pre-Process"
    },
    "Pre-Process": {
      "Type": "Pass",
      "Parameters": {
        "OrderedArray": [
          {
            "OriginalData": {
              "Input.$": "$",
              "ShouldFilterData": false
            }
          },
          {
            "ValuesToCheck": {
              "ListBranchErrors.$": "$.ParralelOutput[?(@.BranchError==true)].BranchError",
              "BranchFailures": true
            }
          },
          {
            "DefaultAlwaysFalse": {
              "ShouldFilterData": false,
              "BranchFailures": false
            }
          }
        ]
      },
      "OutputPath": "$..[?(@.ShouldFilterData == false || @.ListBranchErrors[0] == true)]",
      "Next": "ChoiceStateX"
    },
    "ChoiceStateX": {
      "Type": "Choice",
      "OutputPath": "$.[0].Input",
      "Choices": [
        {
          "Variable": "$[1].BranchFailures",
          "BooleanEquals": true,
          "Next": "NotifyFailure"
        },
        {
          "Variable": "$[1].BranchFailures",
          "BooleanEquals": false,
          "Next": "NotifySuccess"
        }
      ],
      "Default": "NotifyFailure"
    },
    "NotifyFailure": {
      "Type": "Pass",
      "End": true
    },
    "NotifySuccess": {
      "Type": "Pass",
      "Result": "This is a fallback from a task success",
      "End": true
    }
  }
}

【讨论】:

  • 我正在使用类似的方法为并行状态下的每个任务单独处理错误。我想知道 Step Functions 中是否有一种方法可以在“选择”状态下循环遍历结果并检查任务状态而无需在“或”子句中进行硬编码?!这样,如果我在并行状态中添加或删除任务,则不需要修改选择状态。
  • 是的,可以使用额外的 Pass State 和一些 JSONPath。我已更新答案以包含更一般的场景,而无需对选择状态进行硬编码。很抱歉没有详细介绍每个步骤,但状态机 JSON 应该足以弄清楚它是如何工作的。如果您有任何问题,请告诉我。
猜你喜欢
  • 1970-01-01
  • 2019-09-07
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-05-15
  • 1970-01-01
相关资源
最近更新 更多