防止工作流重试重复产生业务效果

AWSBeginner
立即练习

介绍

工作程序保存了订单,却在结果到达工作流之前失败。你将保护这次写入,让重试或重复执行保留原始订单,同时让不同订单仍然能够成功。

请先完成 重试失败步骤并处理永久错误、使用条件写入防止重复订单 和 配置并诊断 Lambda 函数。这个独立 VM 提供未受保护的工作程序;保护逻辑由你添加。

相关认证考点

本实验为以下认证考点提供动手练习。

使用稳定键保护业务写入

本步骤中,部署一个工作程序,将内容相同的重复订单视为已完成的写入。

两次执行对应一个业务订单

不同执行名称可以指向同一个业务订单。内容相同的条件写入冲突返回重复结果,而不会替换该订单。

使用 Terminal 旁的 AWS View,将 CLI 查询与本实验的实际资源和结果对比。保留提供的参考数据。

订单 ID 是业务键:即使任务尝试或工作流执行名称不同,它仍标识预期订单。attribute_not_exists(id) 只允许第一次 PutItem。DynamoDB 原子地评估该条件。当重试未能通过条件竞争时,ReturnValuesOnConditionCheckFailure:ALL_OLD 提供已有项目。只有该项目与预期订单相等时,才返回重复结果;内容冲突仍必须失败,不能悄悄替换订单。

提供的起始程序使用无条件写入。用下面完整的受保护处理程序替换它。它使用配置好的 SDK 环境,而不是个人凭证。诊断尝试独立于业务写入,跟踪实际调用。在模拟的 after-write 模式中,第一次尝试在原生写入之后抛出错误,制造重试必须处理的不确定性。summary 模式读取实际保存的项目。

带引号的 here-document 写入字面 Python 文件。已有 Lambda 处理程序仍为 app.handler。

cd /home/labex/project
cat > app.py <<'PYTHON'
import json
import os
import boto3
from botocore.exceptions import ClientError

class TransientOrderError(Exception):
    pass
class InvalidOrder(Exception):
    pass

def handler(event, context):
    print('WORKFLOW '+json.dumps(event,sort_keys=True))
    db=boto3.client('dynamodb')
    if event.get('stage')=='summary':
        item=db.get_item(
            TableName=os.environ['ORDERS_TABLE'],
            Key={'id':{'S':event['id']}},
        )['Item']
        result={'id':event['id'],'total_cents':int(item['total_cents']['N']),'completed':True}
        print('RESULT '+json.dumps(result,sort_keys=True))
        return result
    if event.get('mode')=='permanent':
        raise InvalidOrder('Order cannot be completed')
    quantity=event['quantity']
    if isinstance(quantity,bool) or not isinstance(quantity,int) or not 1<=quantity<=10:
        raise InvalidOrder('Invalid quantity')
    attempt=db.update_item(
        TableName=os.environ['ATTEMPTS_TABLE'],
        Key={'id':{'S':event['id']}},
        UpdateExpression='ADD attempts :one',
        ExpressionAttributeValues={':one':{'N':'1'}},
        ReturnValues='UPDATED_NEW',
    )['Attributes']['attempts']['N']
    if event.get('mode')=='flaky' and attempt=='1':
        raise TransientOrderError('Dependency temporarily unavailable')
    item={'id':{'S':event['id']},'quantity':{'N':str(quantity)},'total_cents':{'N':str(quantity*250+100)}}
    duplicate=False
    try:
        db.put_item(
            TableName=os.environ['ORDERS_TABLE'],
            Item=item,
            ConditionExpression='attribute_not_exists(id)',
            ReturnValuesOnConditionCheckFailure='ALL_OLD',
        )
    except ClientError as error:
        if (
            error.response['Error']['Code']!='ConditionalCheckFailedException'
            or error.response.get('Item')!=item
        ):
            raise
        duplicate=True
    if event.get('mode')=='after-write' and attempt=='1':
        raise TransientOrderError('Result delivery failed after the write')
    result={'id':event['id'],'quantity':quantity,'total_cents':quantity*250+100,'duplicate':duplicate}
    print('RESULT '+json.dumps(result,sort_keys=True))
    return result
PYTHON

try 执行原生条件写入。只有实际的 ConditionalCheckFailedException 且旧项目内容相同,才被视为重复;其他 SDK 错误会再次抛出。临时异常发生在该写入决定之后,因此工作流会实际重试业务项目已存在的任务。

使用 zip -j 打包文件,它会省略目录路径;通过 --zip-file fileb:// 将二进制归档字节上传到提供的工作程序。响应中的 CodeSha256 标识已部署归档。仅部署不能证明重复保护;运行步骤将测试它。

zip -j function.zip app.py
aws lambda update-function-code \
  --function-name labex-ev05-worker \
  --zip-file fileb://function.zip \
  --query CodeSha256 \
  --output text

创建执行之前,运行部署检查。

连接具有有限重试的限定范围工作流

本步骤中,创建独立的工作流角色和状态机,用于重试工作程序写入后的失败。

Step Functions 需要 states 服务信任关系,以及针对精确函数的独立 InvokeFunction 授权。工作程序保留自己的 DynamoDB 权限;工作流角色只调用它。shell 赋值保存返回的标识符。--query 和 --output text 选择可复用值,file:// 读取字面信任 JSON,shell 将工作程序 ARN 插入权限文档。

cd /home/labex/project
WORKER_NAME=labex-ev05-worker
WORKER_ARN=$(aws lambda get-function-configuration \
  --function-name labex-ev05-worker \
  --query FunctionArn \
  --output text)
aws lambda get-function-configuration \
  --function-name labex-ev05-worker \
  --query '{Name:FunctionName,Role:Role,Runtime:Runtime,Timeout:Timeout}'
aws stepfunctions list-state-machines
cat > workflow-trust.json <<'JSON'
{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Effect": "Allow",
      "Principal": {
        "Service": "states.amazonaws.com"
      },
      "Action": "sts:AssumeRole"
    }
  ]
}
JSON
ROLE_ARN=$(aws iam create-role \
  --role-name labex-ev05-workflow-role \
  --assume-role-policy-document file://workflow-trust.json \
  --query Role.Arn \
  --output text)
cat > workflow-invoke.json <<EOF
{
  "Version": "2012-10-17",
  "Statement": [
    {
      "Effect": "Allow",
      "Action": "lambda:InvokeFunction",
      "Resource": "$WORKER_ARN"
    }
  ]
}
EOF
aws iam put-role-policy \
  --role-name labex-ev05-workflow-role \
  --policy-name InvokeWorker \
  --policy-document file://workflow-invoke.json
aws iam get-role-policy --role-name labex-ev05-workflow-role --policy-name InvokeWorker

空状态机列表确认没有复用之前 VM 的工作流。授权使用函数 ARN,每个 Task 则使用普通函数名称。Retry 只处理 TransientOrderError,间隔为一秒、退避倍数为二,第一次尝试后最多重试两次;Catch 将 InvalidOrder 路由到明确的 Fail。ResultSelector 保留实际 Payload,ResultPath 将其保存在 saved 下,OutputPath 返回实际摘要。

带引号的 here-document 写入字面 ASL。jq --arg 在创建状态机之前替换工作程序占位符。

cat > workflow-template.json <<'JSON'
{
  "StartAt": "CheckQuantity",
  "States": {
    "CheckQuantity": {
      "Type": "Choice",
      "Choices": [
        {
          "Variable": "$.quantity",
          "NumericGreaterThan": 0,
          "Next": "StoreOrder"
        }
      ],
      "Default": "Rejected"
    },
    "StoreOrder": {
      "Type": "Task",
      "Resource": "arn:aws:states:::lambda:invoke",
      "Parameters": {
        "FunctionName": "WORKER_NAME",
        "Payload.$": "$"
      },
      "ResultSelector": {
        "result.$": "$.Payload"
      },
      "ResultPath": "$.saved",
      "Retry": [
        {
          "ErrorEquals": [
            "TransientOrderError"
          ],
          "IntervalSeconds": 1,
          "BackoffRate": 2,
          "MaxAttempts": 2
        }
      ],
      "Catch": [
        {
          "ErrorEquals": [
            "InvalidOrder"
          ],
          "Next": "Rejected",
          "ResultPath": "$.failure"
        }
      ],
      "Next": "ReadSummary"
    },
    "ReadSummary": {
      "Type": "Task",
      "Resource": "arn:aws:states:::lambda:invoke",
      "Parameters": {
        "FunctionName": "WORKER_NAME",
        "Payload": {
          "stage": "summary",
          "id.$": "$.saved.result.id"
        }
      },
      "OutputPath": "$.Payload",
      "End": true
    },
    "Rejected": {
      "Type": "Fail",
      "Error": "OrderRejected",
      "Cause": "Order could not be completed"
    }
  }
}
JSON
jq --arg worker "$WORKER_NAME" '.States.StoreOrder.Parameters.FunctionName=$worker | .States.ReadSummary.Parameters.FunctionName=$worker' workflow-template.json > workflow.json
MACHINE_ARN=$(aws stepfunctions create-state-machine \
  --name labex-ev05-orders \
  --type STANDARD \
  --role-arn "$ROLE_ARN" \
  --definition file://workflow.json \
  --query stateMachineArn \
  --output text)
aws stepfunctions describe-state-machine \
  --state-machine-arn "$MACHINE_ARN" \
  --query '{Name:name,Definition:definition}'

原生定义列出选择性的 Retry/Catch 和两个实际任务。AWS View 显示状态机,尚无执行或订单。运行工作流配置检查。

证明重试与重新执行只产生一次业务写入

本步骤中,运行写入后失败,重复同一订单,再创建不同订单。

执行名称标识工作流运行;订单 ID 标识业务效果。第一次运行使用 after-write 模式。有次数上限的循环读取状态,直到不再为 RUNNING;$(...) 捕获输出,break 退出循环。如果循环后仍为 RUNNING,请检查执行。

FIRST_ARN=$(aws stepfunctions start-execution \
  --state-machine-arn "$MACHINE_ARN" \
  --name after-write-order \
  --input '{"id":"protected-order","quantity":4,"mode":"after-write"}' \
  --query executionArn \
  --output text)
for attempt in $(seq 1 60); do
  STATUS=$(aws stepfunctions describe-execution \
    --execution-arn "$FIRST_ARN" \
    --query status \
    --output text)
  if test "$STATUS" != RUNNING; then break; fi
  sleep 2
done
aws stepfunctions describe-execution \
  --execution-arn "$FIRST_ARN" \
  --query '{Status:status,Output:output}'
aws stepfunctions get-execution-history \
  --execution-arn "$FIRST_ARN" \
  --query 'events[?type==`TaskFailed` || type==`TaskSucceeded`].{Type:type,Error:taskFailedEventDetails.error,Output:taskSucceededEventDetails.output}'
aws dynamodb get-item \
  --table-name labex-ev05-orders \
  --key '{"id":{"S":"protected-order"}}' \
  --query Item

第一次 TaskFailed 是实际写入后的 TransientOrderError。成功的重试返回 duplicate:true;摘要完成,总金额为 1100。唯一的 protected-order 项目数量为 4、总金额为 1100。任务重试成功,没有对该业务键进行第二次被接受的 PutItem。

使用相同订单 ID 和值,以新的执行名称启动执行:

REPEAT_ARN=$(aws stepfunctions start-execution \
  --state-machine-arn "$MACHINE_ARN" \
  --name repeated-order \
  --input '{"id":"protected-order","quantity":4,"mode":"after-write"}' \
  --query executionArn \
  --output text)
for attempt in $(seq 1 60); do
  STATUS=$(aws stepfunctions describe-execution \
    --execution-arn "$REPEAT_ARN" \
    --query status \
    --output text)
  if test "$STATUS" != RUNNING; then break; fi
  sleep 2
done
aws stepfunctions describe-execution \
  --execution-arn "$REPEAT_ARN" \
  --query '{Status:status,Output:output}'
aws stepfunctions get-execution-history \
  --execution-arn "$REPEAT_ARN" \
  --query 'events[?type==`TaskSucceeded`].taskSucceededEventDetails.output'

这次新运行也返回重复存储结果,并读取原始的 1100 摘要。新的执行名称不会让订单变成新业务请求。实际条件失败保留原始项目,而不是覆盖它。

使用不同订单 ID,证明保护逻辑不会拒绝无关任务:

NEW_ARN=$(aws stepfunctions start-execution \
  --state-machine-arn "$MACHINE_ARN" \
  --name different-order \
  --input '{"id":"different-order","quantity":1,"mode":"normal"}' \
  --query executionArn \
  --output text)
for attempt in $(seq 1 60); do
  STATUS=$(aws stepfunctions describe-execution \
    --execution-arn "$NEW_ARN" \
    --query status \
    --output text)
  if test "$STATUS" != RUNNING; then break; fi
  sleep 2
done
aws stepfunctions describe-execution \
  --execution-arn "$NEW_ARN" \
  --query '{Status:status,Output:output}'
aws dynamodb scan --table-name labex-ev05-orders --query Items
aws logs filter-log-events \
  --log-group-name /aws/lambda/labex-ev05-worker \
  --query 'events[].message'

不同订单成功完成,结果为 1/350。虽然存储尝试和摘要读取总共发生了七次实际工作程序调用,但只有两个业务项目。日志和原生任务结果显示原始键的重复判断。AWS View 显示三个成功摘要和两个保留的订单。

下面的示例展示重试和重新执行的实际摘要,以及保留的原始订单和独立的新订单。

AWS View 显示保留的原始订单和不同业务订单 这是针对所测试的相同请求的幂等写入,并不保证任务只运行一次。

运行业务保护检查。

移除工作流资源与模拟结果

本步骤中,删除你已完成执行的状态机、工作流角色、订单和日志,同时保留提供的配套资源。

三次执行都已结束。删除状态机会将它从活动状态机列表中移除。删除角色之前,先移除自己创建的角色策略,然后移除执行所创建的模拟订单和工作程序日志组。

aws stepfunctions delete-state-machine --state-machine-arn "$MACHINE_ARN"
aws iam delete-role-policy --role-name labex-ev05-workflow-role --policy-name InvokeWorker
aws iam delete-role --role-name labex-ev05-workflow-role
aws dynamodb delete-item \
  --table-name labex-ev05-orders \
  --key '{"id":{"S":"protected-order"}}'
aws dynamodb delete-item \
  --table-name labex-ev05-attempts \
  --key '{"id":{"S":"protected-order"}}'
aws dynamodb delete-item \
  --table-name labex-ev05-orders \
  --key '{"id":{"S":"different-order"}}'
aws dynamodb delete-item \
  --table-name labex-ev05-attempts \
  --key '{"id":{"S":"different-order"}}'
aws logs delete-log-group --log-group-name /aws/lambda/labex-ev05-worker

成功读取资源清单,证明剩余资源状态:

aws stepfunctions list-state-machines
aws iam list-roles --query 'Roles[].RoleName'
aws dynamodb scan --table-name labex-ev05-orders --query Items
aws logs describe-log-groups --query logGroups
aws dynamodb scan --table-name labex-ev05-reference --query Items

没有活动状态机、订单或日志组。只保留提供的工作程序角色,参考项目未改变。保留提供的工作程序和表:它们属于环境准备资源,与你创建的工作流和资源不同。网络或身份验证错误绝不能证明删除成功。

移除本实验创建的普通文件:

rm -f app.py function.zip workflow-trust.json workflow-invoke.json workflow-template.json workflow.json

AWS View 显示状态机、执行和订单为空,参考资源仍保留。结束 VM 之前,运行清理检查。

总结

你使用稳定订单键保护原生业务写入,只将内容相同的条件冲突视为重复,并测试了首次写入之后的实际失败。重试和新的工作流执行保留了原始订单;不同业务键创建了自己的结果。你在移除自己创建的工作流资源和模拟结果之前,验证了历史和原生数据。