はじめに
ワーカーは注文を保存しましたが、結果がワークフローに届く前に失敗します。別の注文は引き続き成功できる状態に保ち、再試行や繰り返し実行が元の注文を保持するように、その書き込みを保護します。
先に 失敗したステップを再試行して永続的なエラーを処理する、条件付き書き込みで注文の重複を防ぐ、Lambda 関数を設定して診断する を完了してください。この独立した VM は、保護されていないワーカーを提供します。保護は自分で追加します。
認定試験との関連
このラボでは、次の試験トピックに関連する実践的な演習を行います。
- Solutions Architect – Associate (SAA-C03) · タスク 2.1: 繰り返しても安全なワークフローと保存済み結果の維持。
- Developer – Associate (DVA-C02) · タスク 1.1: 繰り返しても安全なワークフローと保存済み結果の維持。
- DevOps Engineer – Professional (DOP-C02) · タスク 5.1: 基礎演習:繰り返しても安全なワークフローと保存済み結果の維持。
- Solutions Architect – Professional (SAP-C02) · タスク 2.4: 基礎演習:繰り返しても安全なワークフローと保存済み結果の維持。
安定したキーで業務書き込みを保護する
このステップでは、同じ注文の繰り返しを、すでに完了した書き込みとして扱うワーカーをデプロイします。

異なる実行名が同じ業務注文を参照する場合があります。同一内容の条件付き書き込みの競合は、その注文を置き換えずに重複結果を返します。
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 つの成功したサマリーと、両方の保持された注文が表示されます。
以下の例は、再試行と再実行の実際のサマリーを、保持された元の注文と独立した新しい注文とともに示します。
これは、テストした同一リクエストに対する冪等な書き込みであり、タスクが厳密に 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 を終了する前に、後片付けの検証を実行してください。
まとめ
安定した注文キーでネイティブな業務書き込みを保護し、同一内容の条件競合だけを重複として扱い、最初の書き込み後の実際の失敗をテストしました。再試行と新しいワークフロー実行は元の注文を保持し、別の業務キーは独自の結果を作成しました。所有するワークフローリソースと合成結果を削除する前に、履歴とネイティブなデータを確認しました。



