Lambdaで冪等性を担保する方法として、AWS re:PostにPowertools for AWS Lambdaを使った方法が紹介されています。
参考: Lambda 関数を冪等にする方法を教えてください。
上記の記事を参考にして冪等性担保について学ぶことにします。
使うもの
以下を今回は利用します。
- Powertools for AWS Lambda (Python)
- Amazon SQS
- AWS Lambda
- Amazon DynamoDB
SQSは再配信される可能性があるため、SQSで再配信したとしてもLambdaで同じ処理を再実行しないようにします。
PowertoolsのIdempotencyユーティリティを使うと同じ冪等性キーで処理済みの場合、処理を再実行せずに保存済みの結果を返せます。
今回は処理を識別するため、SQSメッセージにjob_idを用意して、job_idを冪等性キーとして処理状態をDynamoDBへ保存します。
AWS公式資料: Powertools for AWS Lambda (Python) - Idempotency
全体の流れ
以下のような図になります。
sequenceDiagram
participant S as SQS
participant W as Worker Lambda
participant P as Powertools Idempotency
participant D as DynamoDB
participant B as 業務処理
S->>W: ① job_id=job-001を受信
W->>P: ① job_idを抽出
P->>D: ②③ status=INPROGRESSを条件付き保存
P->>B: 業務処理を実行
B-->>P: 成功
P->>D: ④ status=COMPLETEDと処理結果を保存
P-->>W: ⑤ 正常終了
Note over S,D: 同じjob_idをもう一度受信
S->>W: ① job_id=job-001を受信
W->>P: ① 同じjob_idを抽出
P->>D: ② 既存レコードを確認
D-->>P: status=COMPLETEDと処理結果
P-->>W: 保存済みの結果を返す
Note over P,B: 業務処理は実行しない
以下、re:Postの「べき等 Lambda 関数ロジックの例」の内容を理解するために一度整理します(原文を若干変更して書いています)。
① イベントの一意の属性の値を抽出
SQSメッセージのjob_idを今回は一意の属性の値として扱います。
{
"job_id": "job-001"
}
② 条件式を使用してDynamoDBに保存する
条件式とは、DynamoDBのPutItemやUpdateItemで、条件を満たす場合にだけ書き込むための式です。
PowertoolsのIdempotencyユーティリティはこの条件式を利用し、冪等性キーが存在しない場合のみ保存します。
反対に、冪等性キーが既に存在する場合はPowertools側で既存レコードを確認し、処理が完了していれば保存済みの結果を返します。
③ レコード属性に対してアクションを実行
最初の処理では、DynamoDBへstatus=INPROGRESSを保存してから業務処理を実行します。
先に処理中として保存することで、同じjob_idの処理が重複して実行されることを防ぎます。
④ ステータスを更新する
業務処理が成功すると、DynamoDBへstatus=COMPLETEDと処理結果を保存します。
⑤ アクションを完了
DynamoDBへの処理結果の保存まで完了すると、Lambdaは正常終了します。
なおINPROGRESSの間に同じjob_idを受け取った場合、まだ返せる処理結果がないためPowertoolsはIdempotencyAlreadyInProgressErrorを発生させます。
今回はCOMPLETEDになった後の重複だけを確認します。
AWS公式資料: Powertools for AWS Lambda (Python) - Handling concurrent executions with the same payload
実装内容
まず、AWS CDKで冪等性レコードを保存するDynamoDBテーブルとWorker Lambdaを作ります。
主要部分のみ抜粋します。SQSやLogGroupなどは省略しています。
idempotency_table = dynamodb.Table(
self,
"IdempotencyTable",
partition_key=dynamodb.Attribute(
name="id",
type=dynamodb.AttributeType.STRING,
),
billing_mode=dynamodb.BillingMode.PAY_PER_REQUEST,
time_to_live_attribute="expiration",
removal_policy=RemovalPolicy.DESTROY,
)
worker_function = lambda_.Function(
self,
"WorkerFunction",
runtime=lambda_.Runtime.PYTHON_3_13,
handler="handler.lambda_handler",
code=lambda_.Code.from_asset(...),
environment={"IDEMPOTENCY_TABLE": idempotency_table.table_name},
)
idempotency_table.grant_read_write_data(worker_function)
※code=lambda_.Code.from_asset(...)では、handler.pyとPowertoolsなどの依存ライブラリをLambdaへ配置する箇所ですが、割愛するため(...)のように書いています。
persistence_layer = DynamoDBPersistenceLayer(
table_name=os.environ["IDEMPOTENCY_TABLE"],
)
config = IdempotencyConfig(
event_key_jmespath="job_id",
)
@idempotent_function(
data_keyword_argument="job",
persistence_store=persistence_layer,
config=config,
)
def process_job(job: dict[str, str]) -> dict[str, object]:
return run_business_operation(job)
Worker Lambdaでは、SQSメッセージ1件の処理をidempotent_functionで囲みます。
run_business_operationは、業務処理に見立てた関数です。
実行されたことを確認するため、実行ごとにbusiness_operation_executedというログとexecution_idを作ります。
data_keyword_argumentに指定したjob(実態はjob_id)を取り出し、冪等性キーとして使います。
実践内容
同じjob_idのメッセージを2回送りWorker LambdaのログとDynamoDBを確認します。
デプロイ
対象のAWSアカウントとリージョンを確認してからデプロイします。
uv sync --extra dev
aws sts get-caller-identity
aws configure get region
uv run -- npx --yes aws-cdk@2 deploy --outputs-file cdk-outputs.json
--outputs-fileで、AWS CDKの出力値をcdk-outputs.jsonへ保存しています。次のコマンドで、メッセージの送信先、DynamoDBのテーブル名、Worker LambdaのLogGroup名を取得します。
# SQSメッセージの送信先
QUEUE_URL=$(python -c 'import json; print(json.load(open("cdk-outputs.json"))["LambdaIdempotencyStack"]["QueueUrl"])')
# 冪等性レコードを確認するDynamoDBテーブル
TABLE_NAME=$(python -c 'import json; print(json.load(open("cdk-outputs.json"))["LambdaIdempotencyStack"]["IdempotencyTableName"])')
# Worker Lambdaのログ出力先
LOG_GROUP=$(python -c 'import json; print(json.load(open("cdk-outputs.json"))["LambdaIdempotencyStack"]["WorkerLogGroupName"])')
1回目のメッセージを送ります。
aws sqs send-message \
--queue-url "$QUEUE_URL" \
--message-body '{"job_id":"job-001"}'
aws logs tail "$LOG_GROUP" --since 5m --format short
business_operation_executedを確認したら、同じjob_idでもう一度送ります。
aws sqs send-message \
--queue-url "$QUEUE_URL" \
--message-body '{"job_id":"job-001"}'
aws logs tail "$LOG_GROUP" --since 5m --format short
DynamoDBのレコードも確認します。
aws dynamodb scan \
--table-name "$TABLE_NAME" \
--projection-expression "id,#status" \
--expression-attribute-names '{"#status":"status"}'
実測結果
2026年8月13日にap-northeast-1で実行しました。
同じjob_idを持つ、別々のSQSメッセージを2回送りました。
job_id=idempotency-20260813-002
1回目 message_id=e01a4092-37ff-4d12-aaf8-e7976086e641
2回目 message_id=853c0f20-4df0-4134-ab6e-0fe46e9fe056
message_idは、SQSがメッセージごとに付ける識別子です。
AWS公式資料: Amazon SQSメッセージイベントの例
Worker Lambdaのログは次のようになりました。
1回目 event=business_operation_executed job_id=idempotency-20260813-002 execution_id=dde5996a-810a-47b3-b952-ba1d43407326
2回目 event=job_skipped job_id=idempotency-20260813-002 sqs_message_id=853c0f20-4df0-4134-ab6e-0fe46e9fe056
1回目はbusiness_operation_executedが出ました。
2回目はjob_skippedが出て、business_operation_executedは増えませんでした(今回省略したログ出力の処理の中に、スキップする処理を仕込んでいました)。
DynamoDBには、1回目の処理結果が保存されていました。
status=COMPLETED
job_id=idempotency-20260813-002
execution_id=dde5996a-810a-47b3-b952-ba1d43407326
同じjob_idを持つ別々のSQSメッセージを2回送っても、業務処理は1回だけ実行されることを確認できました。