ワークフローの再試行で業務上の効果を繰り返さないようにする

AWSBeginner
オンラインで実践に進む

はじめに

ワーカーは注文を保存しましたが、結果がワークフローに届く前に失敗します。別の注文は引き続き成功できる状態に保ち、再試行や繰り返し実行が元の注文を保持するように、その書き込みを保護します。

先に 失敗したステップを再試行して永続的なエラーを処理する、条件付き書き込みで注文の重複を防ぐ、Lambda 関数を設定して診断する を完了してください。この独立した VM は、保護されていないワーカーを提供します。保護は自分で追加します。

認定試験との関連

このラボでは、次の試験トピックに関連する実践的な演習を行います。

安定したキーで業務書き込みを保護する

このステップでは、同じ注文の繰り返しを、すでに完了した書き込みとして扱うワーカーをデプロイします。

2 つの実行と 1 つの業務注文

異なる実行名が同じ業務注文を参照する場合があります。同一内容の条件付き書き込みの競合は、その注文を置き換えずに重複結果を返します。

Terminal の隣で AWS View を使い、CLI クエリをこのラボの実際のリソースと結果と比較します。提供された参照データを保持してください。

注文 ID は業務キーです。タスクの試行やワークフロー実行名が異なっても、意図する注文を識別します。attribute_not_exists(id) は最初の PutItem だけを許可します。DynamoDB は、その条件をアトミックに評価します。ReturnValuesOnConditionCheckFailure:ALL_OLD は、再試行が条件の競合に負けた場合に既存のアイテムを提供します。その既存アイテムが意図する注文と等しい場合だけ、重複結果を返します。競合する内容は、黙って注文を置き換えず、引き続き失敗させる必要があります。

提供された初期コードには、無条件の書き込みがあります。以下の完全な保護付きハンドラーに置き換えます。個人の認証情報ではなく、設定済みの SDK 環境を使います。診断用試行記録は、業務書き込みとは別に実際の呼び出しを追跡します。合成された after-write モードでは、最初の試行がネイティブな書き込み後にエラーを発生させ、再試行で扱う必要がある不確実性を生みます。summary モードは、実際の保存アイテムを読み取ります。

引用符付きのヒアドキュメントは、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 権限を保持し、ワークフローロールは呼び出しだけを行います。シェルの代入は返された識別子を保存します。--query と --output text は再利用可能な値を選択し、file:// は信頼 JSON をそのまま読み込みます。シェルは、ワーカー 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 だけを扱い、1 秒の間隔と倍増するバックオフで、最初の試行後に最大 2 回再試行します。Catch は InvalidOrder を明示的な Fail に振り分けます。ResultSelector は実際の Payload を保持し、ResultPath はそれを saved の下に保存し、OutputPath は実際のサマリーを返します。

引用符付きのヒアドキュメントは 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 と 2 つの実際のタスクがあります。AWS View には、実行や注文がまだないステートマシンが表示されます。ワークフロー設定の検証を行ってください。

再試行と再実行を通じた 1 回の業務書き込みを確認する

このステップでは、書き込み後の失敗を実行し、同じ注文を繰り返して、別の注文を作成します。

実行名はワークフローの実行を識別し、注文 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 で完了します。1 件の protected-order アイテムは数量 4、合計 1100 です。その業務キーに対する 2 回目の受け入れられた 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 で成功します。保存試行とサマリー読み取りを通じて、実際のワーカー呼び出しは 7 回発生しましたが、業務アイテムは 2 件です。ログとネイティブなタスク結果は、元のキーへの重複判定を示します。AWS View には、3 つの成功したサマリーと、両方の保持された注文が表示されます。

以下の例は、再試行と再実行の実際のサマリーを、保持された元の注文と独立した新しい注文とともに示します。

AWS View に保持された元の注文と別の業務注文が表示される これは、テストした同一リクエストに対する冪等な書き込みであり、タスクが厳密に 1 回だけ実行されるという約束ではありません。

業務保護の検証を実行してください。

ワークフローリソースと合成結果を削除する

このステップでは、提供された準備リソースを保持しながら、自分の完了したステートマシン、ワークフローロール、注文、ログを削除します。

3 つの実行はすべて終了しています。ステートマシンの削除は、アクティブなステートマシン一覧からそれを取り除きます。ロールを削除する前に所有するロールポリシーを削除し、その後、合成注文と実行によって作成されたワーカーロググループを削除します。

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 を終了する前に、後片付けの検証を実行してください。

まとめ

安定した注文キーでネイティブな業務書き込みを保護し、同一内容の条件競合だけを重複として扱い、最初の書き込み後の実際の失敗をテストしました。再試行と新しいワークフロー実行は元の注文を保持し、別の業務キーは独自の結果を作成しました。所有するワークフローリソースと合成結果を削除する前に、履歴とネイティブなデータを確認しました。