重複した効果を生まずにキュージョブを処理する

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

はじめに

同じ注文の繰り返されるジョブは、完了した業務結果を保持する必要があります。提供されたコンシューマーにアトミックな書き込み保護を追加し、繰り返される注文と独立した注文をテストして、処理済みとしての確認が成功することを確かめます。

先に デッドレターキューで失敗したジョブを隔離する、条件付き書き込みで注文の重複を防ぐ、Lambda 関数を設定して診断する を完了してください。この独立した VM は、初期状態のワーカー、ロール、空の注文テーブルを提供します。

認定試験との関連

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

空のキューをワーカーに接続する

このステップでは、業務ジョブを送る前にキューを作成し、提供されたコンシューマーを接続します。

Terminal の隣で AWS View を使い、現在のキュー、ワーカーの結果、保存された注文を比較します。無関係な参照データを保持してください。

cd /home/labex/project

Standard キューを作成します。コマンド置換 $(...) は、後の操作のために、返されたキュー URL を保存します。

QUEUE_URL=$(aws sqs create-queue --queue-name labex-q05-jobs --attributes VisibilityTimeout=30 --query QueueUrl --output text)

ARN を選択し、サイズ 1 のイベントソースマッピングを作成します。

QUEUE_ARN=$(aws sqs get-queue-attributes --queue-url "$QUEUE_URL" --attribute-names QueueArn --query Attributes.QueueArn --output text)
MAPPING_ID=$(aws lambda create-event-source-mapping --function-name labex-q05-worker --event-source-arn "$QUEUE_ARN" --batch-size 1 --enabled --query UUID --output text)

接続を確認します。

aws lambda get-event-source-mapping --uuid "$MAPPING_ID" --query '{Source:EventSourceArn,Function:FunctionArn,Batch:BatchSize,State:State}'

送信元は labex-q05-jobs、関数は labex-q05-worker、バッチサイズは 1、状態は Enabled となるはずです。提供されたロールは、このキューだけをポーリング・削除し、注文テーブルだけに書き込み、実行ログも記録できます。AWS View には空のキューが表示され、注文も書き込み判定もありません。次のステップで保護をデプロイするまで、ジョブを送らないでください。

業務キーに対するアトミックな保護をデプロイする

このステップでは、ワーカーの無条件書き込みを置き換え、最初の保護されたジョブが成功することを確認します。

メッセージ ID と業務キー

異なるキューメッセージに、同じ業務 id が含まれる場合があります。ワーカーは、注文を再び書き込まずに、一致する重複を処理済みとして確認します。

冪等なコンシューマーは、繰り返される作業が到着したときに、同じ業務結果を保持します。DynamoDB は、PutItem と同じ書き込み操作で ConditionExpression='attribute_not_exists(id)' を評価します。これにより、別々の読み取りと書き込みの競合を避けられます。キーは業務の id です。SQS の messageId を使うと、プロデューサーが同じ注文の新しいコピーを送ったときに、再び書き込めてしまいます。

条件が失敗すると、ReturnValuesOnConditionCheckFailure='ALL_OLD' が既存のアイテムを返します。コードは ConditionalCheckFailedException だけを捕捉し、そのアイテムを要求されたアイテムと比較します。同一の詳細なら完了済みの重複です。競合する詳細や別のエラーは、黙って処理済みとして確認せず、引き続き失敗させます。

以下の完全なワーカーを書き込みます。cat > app.py はファイルを置き換え、ヒアドキュメントは終端の PY までの内容を提供します。'PY' を引用符で囲むことで、Python コード内のシェル展開を防ぎます。

cat > app.py <<'PY'
import json
import os
import boto3
from botocore.exceptions import ClientError


def handler(event, context):
    print('EVENT '+json.dumps(event,sort_keys=True))
    if set(event)!={'Records'} or len(event['Records'])!=1:
        raise ValueError('Expected one SQS record')
    job=json.loads(event['Records'][0]['body'])
    print('JOB '+json.dumps(job,sort_keys=True))
    if set(job)!={'id','quantity'} or not isinstance(job['id'],str):
        raise ValueError('Use an id and quantity')
    quantity=job['quantity']
    if isinstance(quantity,bool) or not isinstance(quantity,int) or not 1<=quantity<=10:
        raise ValueError('Quantity must be an integer from 1 to 10')
    database=boto3.client('dynamodb',endpoint_url='http://127.0.0.1:5000',region_name='us-east-1')
    item={'id':{'S':job['id']},'quantity':{'N':str(quantity)},'total_cents':{'N':str(quantity*250+100)}}
    try:
        database.put_item(TableName=os.environ['TABLE_NAME'],Item=item,
                          ConditionExpression='attribute_not_exists(id)',
                          ReturnValuesOnConditionCheckFailure='ALL_OLD')
        result={'id':job['id'],'quantity':quantity,'processed':True,'duplicate':False}
    except ClientError as error:
        if error.response['Error']['Code']!='ConditionalCheckFailedException':
            raise
        if error.response.get('Item')!=item:
            raise ValueError('Order ID already exists with different details') from error
        result={'id':job['id'],'quantity':quantity,'processed':False,'duplicate':True}
    print('RESULT '+json.dumps(result,sort_keys=True))
    return result
PY

イベント解析と数量チェックは、提供された前提部分です。主な変更は、単一の条件付き put_item、その限定された失敗処理、processed/duplicate の結果です。重複は成功として返るため、業務結果がすでに保持された後、コンシューマーがその SQS コピーを処理済みとして確認できます。

ZIP のルートに app.py をパッケージ化します。既存関数のハンドラーは app.handler なので、モジュールのファイル名が重要です。

zip -q function.zip app.py

fileb:// でバイナリ ZIP をアップロードします。

aws lambda update-function-code --function-name labex-q05-worker --zip-file fileb://function.zip --query CodeSha256 --output text

返されたハッシュは、デプロイされたコードを識別します。実際の処理がその動作を証明します。最初の注文を送ります。

aws sqs send-message --queue-url "$QUEUE_URL" --message-body '{"id":"dedup-order","quantity":4}'

キューが空になり、ワーカーが processed: true を報告して、注文が表示されるまで AWS View を観察します。その後、注文を読み取ります。

aws dynamodb get-item --table-name labex-q05-orders --key '{"id":{"S":"dedup-order"}}' --consistent-read --query Item

数量 4 と合計 1100 が返るはずです。書き込み判定には、条件、Accepted、以前のアイテムがないこと、書き込み後のこの保存アイテムが表示されます。条件の設定やアップロード済みファイルだけでは、作業成功を証明できません。

書き込みを繰り返さずに繰り返されるジョブを消費する

このステップでは、同じ業務注文の新しいコピーを 2 件送り、別の注文も独立して完了することを確認します。

同じ業務ペイロードを 2 回送ります。各 SendMessage レスポンスには新しい SQS メッセージ ID がありますが、注文キーは dedup-order のままです。

aws sqs send-message --queue-url "$QUEUE_URL" --message-body '{"id":"dedup-order","quantity":4}'
aws sqs send-message --queue-url "$QUEUE_URL" --message-body '{"id":"dedup-order","quantity":4}'

独立した業務キーを送ります。

aws sqs send-message --queue-url "$QUEUE_URL" --message-body '{"id":"other-order","quantity":1}'

4 回の実際の実行が返り、キューが空になるまで AWS View を観察します。2 回の繰り返し実行は、processed: false、duplicate: true を報告します。DynamoDB は両方の条件付き書き込みを ConditionalCheckFailedException で拒否し、各アイテムは書き込み前後で変わりません。最初の注文と other-order には、それぞれ 1 回の受け入れられた書き込みがあります。コンシューマーは、すでに完了した作業を再試行する代わりに、繰り返しのコピーを処理済みとして確認します。

業務結果を読み取ります。

aws dynamodb scan --table-name labex-q05-orders --query Items

dedup-order/4/1100 と other-order/1/350 の 2 件だけがあるはずです。アイテムの順序は異なる場合があります。実際のワーカー結果を確認します。ログイベントには複数行が含まれる場合があります。このパイプラインは、ネイティブな JSON 出力を jq -r に送り、各イベントを行に分割して、RESULT 行だけを選択します。表示結果にレシートハンドルは含まれません。

aws logs filter-log-events --log-group-name /aws/lambda/labex-q05-worker --filter-pattern '"RESULT"' --output json | jq -r '.events[].message | split("\n")[] | select(startswith("RESULT "))'

処理済みとしての確認をチェックします。

aws sqs get-queue-attributes --queue-url "$QUEUE_URL" --attribute-names ApproximateNumberOfMessages ApproximateNumberOfMessagesNotVisible

両方の件数は 0 です。実際のメッセージ配信と関数実行は 4 回、受け入れられた業務書き込みは 2 回、ネイティブな条件拒否は 2 回ありました。テーブルのアイテム数だけを数えても、無条件の上書きは分かりません。実際の判定と対応するコンシューマー結果も使ってください。

この保護は 1 つの業務アイテムを保護します。DynamoDB 書き込みと SQS の処理済み確認を、単一のアトミックな操作にするものではありません。注文レコードが永続的な重複保護です。それを削除すると、同じキーが再びアイテムを作成できるため、保持期間と業務キーの設計は実際のシステムの方針の一部です。これらの合成された繰り返しは、同じ業務作業を安全に扱えることを示します。本番環境の並行処理やエンドツーエンドで厳密に 1 回の処理を保証するものではありません。

AWS View の例:4 回の実際の実行で 2 件の注文と、既存アイテムを保持する 2 回の条件拒否が発生し、キューは空になる

コンシューマー接続と業務結果を削除する

このステップでは、重複処理を確認した後、所有する接続とリソースを削除します。

送信元キューを削除する前に、自分のイベントソースマッピングを停止します。

aws lambda delete-event-source-mapping --uuid "$MAPPING_ID" --query UUID --output text
aws sqs delete-queue --queue-url "$QUEUE_URL"

両方の合成業務アイテムと、このラボの実行ログを削除します。

aws dynamodb delete-item --table-name labex-q05-orders --key '{"id":{"S":"dedup-order"}}'
aws dynamodb delete-item --table-name labex-q05-orders --key '{"id":{"S":"other-order"}}'
aws logs delete-log-group --log-group-name /aws/lambda/labex-q05-worker

不在を確認し、無関係な参照データを保持します。

aws lambda list-event-source-mappings --function-name labex-q05-worker --query EventSourceMappings
aws sqs list-queues
aws dynamodb scan --table-name labex-q05-orders --query Items
aws dynamodb scan --table-name labex-q05-reference --query Items

マッピングと注文は空になり、キュー URL は残らず、参照アイテムには引き続き keep unchanged があります。提供された関数とテーブル構造は、この環境に残ります。ネットワークエラーや認証エラーは削除の証拠になりません。AWS View は過去の書き込み判定を保持しながら、空の現在のリソースを表示します。

通常のコードファイルとアーカイブファイルを削除します。

rm -f app.py function.zip

環境を終了する前に、後片付けの検証を実行してください。

まとめ

実際のキューコンシューマーに、DynamoDB の業務キーに対するアトミックな条件をデプロイしました。最初の注文は完了しました。別々に送った 2 件のコピーは、ネイティブな条件拒否が元のアイテムを保持した後に、消費され、処理済みとして確認されました。別のキーは独自の結果を生成できました。対象を限定した後片付けの前に、実際の書き込み、ワーカー結果、空のキュー状態で結果を確認しました。

チャレンジでは回数を限定した失敗隔離を適用し、その後のサーバーレスプロジェクトで、DLQ 復旧とこの業務キー保護を組み合わせます。