SNS と SQS で通知をファンアウトする

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

はじめに

出荷処理にはすべての注文通知が必要ですが、分析には新しい注文だけが必要です。2 つのキューに向けて 1 回発行し、配信権限を設定して、分析用通知をフィルタリングします。

先に SQS でジョブを送信して消費する と リソースポリシーでバケットを保護する を完了してください。この独立した VM は、設定済み CLI と無関係な参照データを提供します。キューは通知の移動先です。このユニットに Lambda コンシューマーは必要ありません。

認定試験との関連

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

トピックを作成してキューへの配信を許可する

このステップでは、トピックと 2 つの空のキューを作成し、そのトピックだけが配信できるように権限を設定します。

Amazon Simple Notification Service(SNS) は、トピックのサブスクリプションを通じて、1 回の発行を配信します。Terminal の隣で AWS View を使い、トピックの接続、キュー本文、フィルターを比較してください。参照データを保持します。

cd /home/labex/project

プロデューサーの SNS トピックを作成します。コマンド置換 $(...) は、後続のコマンドのために、返された ARN をシェル変数に保存します。

TOPIC_ARN=$(aws sns create-topic --name labex-q04-orders --query TopicArn --output text)

独立した移動先を作成し、キュー URL を保存します。

FULFILLMENT_URL=$(aws sqs create-queue --queue-name labex-q04-fulfillment --query QueueUrl --output text)
ANALYTICS_URL=$(aws sqs create-queue --queue-name labex-q04-analytics --query QueueUrl --output text)

サブスクリプションは、キュー ARN を配信先に使います。両方の識別子を選択します。

FULFILLMENT_ARN=$(aws sqs get-queue-attributes --queue-url "$FULFILLMENT_URL" --attribute-names QueueArn --query Attributes.QueueArn --output text)
ANALYTICS_ARN=$(aws sqs get-queue-attributes --queue-url "$ANALYTICS_URL" --attribute-names QueueArn --query Attributes.QueueArn --output text)

キューポリシーは、SNS サービスにこの正確なキューへの sqs:SendMessage を許可し、aws:SourceArn を自分のトピックに限定します。サブスクリプション設定だけでは、配信権限は付与されません。キューごとに通常の JSON ポリシーを書き込みます。シェルは ARN 変数を挿入します。

cat > fulfillment-policy.json <<EOF
{
  "Version": "2012-10-17",
  "Statement": [{
    "Effect": "Allow",
    "Principal": {"Service": "sns.amazonaws.com"},
    "Action": "sqs:SendMessage",
    "Resource": "$FULFILLMENT_ARN",
    "Condition": {"ArnEquals": {"aws:SourceArn": "$TOPIC_ARN"}}
  }]
}
EOF
cat > analytics-policy.json <<EOF
{
  "Version": "2012-10-17",
  "Statement": [{
    "Effect": "Allow",
    "Principal": {"Service": "sns.amazonaws.com"},
    "Action": "sqs:SendMessage",
    "Resource": "$ANALYTICS_ARN",
    "Condition": {"ArnEquals": {"aws:SourceArn": "$TOPIC_ARN"}}
  }]
}
EOF

SQS の Policy 属性は JSON 文字列を保持します。--rawfile はポリシーファイルをその文字列に読み込み、その後 CLI が通常の属性ファイルを使って設定します。

jq -n --rawfile policy fulfillment-policy.json '{Policy:$policy}' > fulfillment-attributes.json
jq -n --rawfile policy analytics-policy.json '{Policy:$policy}' > analytics-attributes.json
aws sqs set-queue-attributes --queue-url "$FULFILLMENT_URL" --attributes file://fulfillment-attributes.json
aws sqs set-queue-attributes --queue-url "$ANALYTICS_URL" --attributes file://analytics-attributes.json

設定された権限を読み取ります。

aws sqs get-queue-attributes --queue-url "$FULFILLMENT_URL" --attribute-names QueueArn Policy
aws sqs get-queue-attributes --queue-url "$ANALYTICS_URL" --attribute-names QueueArn Policy

両方のポリシーは、自身のキューと同じ注文トピックを指定しています。AWS View には、トピック 1 つ、空のキュー 2 つが表示され、サブスクリプションはまだありません。

両方のキューを購読させて分析用通知をフィルタリングする

このステップでは、各キューをトピックに接続し、分析が受け取る通知を選択します。

SNS ファンアウトとサブスクリプションフィルター

分析用サブスクリプションは、メッセージ属性 kind=created を選択します。各キューは独自のコピーを受け取ります。

サブスクリプションは、プロトコルとエンドポイントを SNS トピックへ接続します。sqs の場合、エンドポイントはキュー ARN です。各接続を設定し、後で削除できるように、サブスクリプション ARN を保存します。

FULFILLMENT_SUB=$(aws sns subscribe --topic-arn "$TOPIC_ARN" --protocol sqs --notification-endpoint "$FULFILLMENT_ARN" --query SubscriptionArn --output text)
ANALYTICS_SUB=$(aws sns subscribe --topic-arn "$TOPIC_ARN" --protocol sqs --notification-endpoint "$ANALYTICS_ARN" --query SubscriptionArn --output text)

フィルターポリシーは、1 つのサブスクリプションへの通知を選択します。デフォルトのフィルター範囲はメッセージ属性です。分析は値が created の String 属性 kind だけを受け入れ、出荷処理にはフィルターがありません。

aws sns set-subscription-attributes --subscription-arn "$ANALYTICS_SUB" --attribute-name FilterPolicy --attribute-value '{"kind":["created"]}'

両方のサブスクリプションと分析用属性を確認します。

aws sns list-subscriptions-by-topic --topic-arn "$TOPIC_ARN"
aws sns get-subscription-attributes --subscription-arn "$ANALYTICS_SUB"

2 つの SQS エンドポイントと分析用フィルターが返るはずです。Raw メッセージ配信はデフォルトで無効なので、キュー本文には、トピック、メッセージ ID、元のメッセージ文字列、属性を含む SNS 通知エンベロープが入ります。発行するまでキューは空のままです。

ファンアウト、フィルタリング、配信境界を確認する

このステップでは、2 つの通知を発行し、実際の独立した配信を確認します。

created イベントを発行します。JSON メッセージは業務ペイロードです。このサブスクリプションフィルターが評価するのは、別に指定する String 属性です。

aws sns publish --topic-arn "$TOPIC_ARN" --message '{"id":"fanout-order","kind":"created"}' --message-attributes '{"kind":{"DataType":"String","StringValue":"created"}}'

同じ注文 ID で updated イベントを発行します。

aws sns publish --topic-arn "$TOPIC_ARN" --message '{"id":"fanout-order","kind":"updated"}' --message-attributes '{"kind":{"DataType":"String","StringValue":"updated"}}'

AWS View を観察します。出荷処理には 2 つの通知があり、分析には created 通知だけがあります。件数を確認します。

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

確認のために可視性 0 で受信します。これにより、合成メッセージを処理中に保持する代わりに、後の後片付けに向けてすぐに再び利用可能にします。この例は、独立したコピーを証明するものであり、処理や処理済みとしての確認ではありません。fromjson は各エンベロープ文字列を読み取り、最後のフィールド選択はここで使う通知フィールドを表示します。

aws sqs receive-message \
  --queue-url "$FULFILLMENT_URL" \
  --max-number-of-messages 10 \
  --visibility-timeout 0 \
  --output json | jq '[.Messages[].Body | fromjson | {Type, TopicArn, MessageId, Message, MessageAttributes}]'
aws sqs receive-message \
  --queue-url "$ANALYTICS_URL" \
  --max-number-of-messages 10 \
  --visibility-timeout 0 \
  --output json | jq '[.Messages[].Body | fromjson | {Type, TopicArn, MessageId, Message, MessageAttributes}]'

本文は SNS エンベロープです。Type は Notification、TopicArn は自分のトピックで、Message は元の JSON 文字列を含みます。created 通知は、両方のキューで同じ SNS MessageId を持ちますが、各キューは独自の SQS メッセージを保持します。出荷処理は、分析がフィルタリングで除外した updated イベントも保持します。独立したコンシューマーは、自分のコピーを個別に処理済みとして確認できます。

次に、トピック権限が必要な理由をテストします。出荷処理の aws:SourceArn を、一時的に無関係な合成トピック ARN へ設定します。

jq --arg wrong 'arn:aws:sns:us-east-1:123456789012:labex-q04-unrelated' '.Statement[0].Condition.ArnEquals["aws:SourceArn"]=$wrong' fulfillment-policy.json > wrong-source-policy.json
jq -n --rawfile policy wrong-source-policy.json '{Policy:$policy}' > wrong-source-attributes.json
aws sqs set-queue-attributes --queue-url "$FULFILLMENT_URL" --attributes file://wrong-source-attributes.json

別の合成注文に対する updated 通知を発行します。SNS は発行を受け付けますが、出荷処理はこのトピックからの配信を認可しなくなり、分析は updated 属性をフィルタリングで除外します。

aws sns publish --topic-arn "$TOPIC_ARN" --message '{"id":"blocked-order","kind":"updated"}' --message-attributes '{"kind":{"DataType":"String","StringValue":"updated"}}'

キュー件数が 2 と 1 のままで、AWS View に blocked-order のエンベロープがないことを確認します。

aws sqs get-queue-attributes --queue-url "$FULFILLMENT_URL" --attribute-names ApproximateNumberOfMessages
aws sqs get-queue-attributes --queue-url "$ANALYTICS_URL" --attribute-names ApproximateNumberOfMessages

発行レスポンスの成功は、SNS が通知を受け付けたことを証明します。移動先への配信には独自の証拠が必要です。このステップの検証前に、正確な元のポリシーを復元します。

aws sqs set-queue-attributes --queue-url "$FULFILLMENT_URL" --attributes file://fulfillment-attributes.json

この作業環境では、同一アカウント内の即時キュー配信と String 属性フィルターをテストします。本番環境の SNS フィルター変更は、反映に最大 15 分かかる場合があります。配信失敗にはサービスの再試行動作もあります。この短い権限実験は、本番環境の再試行モデルではありません。

AWS View の例:出荷処理は created と updated 通知を保持し、分析は同じ SNS メッセージ ID の created 通知だけを保持する

接続と使い捨て通知を削除する

このステップでは、無関係な参照データを保持しながら、両方のサブスクリプション、トピック、キューを削除します。

最初に両方のエンドポイントの購読を解除します。

aws sns unsubscribe --subscription-arn "$FULFILLMENT_SUB"
aws sns unsubscribe --subscription-arn "$ANALYTICS_SUB"

トピックとキューを削除します。キューの削除は合成通知を破棄します。これらの通知について業務処理を行ったとは主張していません。

aws sns delete-topic --topic-arn "$TOPIC_ARN"
aws sqs delete-queue --queue-url "$FULFILLMENT_URL"
aws sqs delete-queue --queue-url "$ANALYTICS_URL"

成功するネイティブなクエリで、後片付けを証明します。

aws sns list-topics
aws sns list-subscriptions
aws sqs list-queues
aws dynamodb scan --table-name labex-q04-reference --query Items

トピックとサブスクリプションの一覧は空になり、キュー URL は存在せず、参照アイテムには引き続き keep unchanged があります。認証やネットワークの失敗は削除の証拠になりません。AWS View にも同じ空のリソースが表示されます。

通常の設定ファイルを削除します。

rm -f fulfillment-policy.json analytics-policy.json fulfillment-attributes.json analytics-attributes.json wrong-source-policy.json wrong-source-attributes.json

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

まとめ

1 つの SNS トピックを 2 つの独立した SQS キューに接続し、配信を正確な送信元トピックに限定して、メッセージ属性で分析用通知をフィルタリングしました。実際の通知エンベロープで、共通の発行イベントと個別のキューコピーを確認しました。また、SNS が発行を受け付けることと、配信が移動先に到達することを区別し、その後、参照データを保持しながら自分のリソースを削除しました。

キューのコンシューマーは、引き続き繰り返し配信に対応する必要があります。次のユニットでは、DynamoDB のアトミックな条件を適用して、繰り返されるジョブが重複した業務上の効果を生まないようにします。