使用 SNS 和 SQS 扇出通知

AWSBeginner
立即练习

介绍

履约系统需要每一条订单通知,而分析系统只需要新订单。你将通过一次发布向两个队列发送通知,配置投递权限,并过滤分析系统的通知。

请先完成 使用 SQS 发送和消费任务 和 使用资源策略保护存储桶。这个独立 VM 提供配置好的 CLI 和无关的参考数据。队列是通知目标;本实验不需要 Lambda 消费者。

创建主题并授予队列投递权限

本步骤中,创建一个主题和两个空队列,只允许该主题向它们投递。

Amazon Simple Notification Service(SNS)通过主题订阅分发一次发布的内容。使用 Terminal 旁的 AWS View 对比主题连接、队列正文和过滤器;保留参考数据。

cd /home/labex/project

创建生产者的 SNS 主题。命令替换 $(...) 将返回的 ARN 保存到 shell 变量中,供后续命令使用:

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 策略;shell 会插入 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 显示一个主题、两个空队列,尚无订阅。

订阅两个队列并过滤分析通知

本步骤中,将每个队列连接到主题,并选择分析系统接收哪些通知。

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)

过滤策略为一个订阅选择通知。默认过滤范围是消息属性。分析系统只接受值为 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"

预期看到两个 SQS 端点和分析过滤器。原始消息投递(Raw message delivery)默认关闭,因此队列正文会包含 SNS 通知封装,其中包括主题、消息 ID、原始消息字符串和属性。发布之前,队列保持为空。

证明扇出、过滤与投递边界

本步骤中,发布两条通知,并检查实际的独立投递。

发布 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:履约队列有两条通知,分析队列只有 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

以零可见性超时领取消息进行检查。这会立即释放这些模拟消息,供稍后清理,而不是让它们保持处理中。此示例证明独立副本的存在,不代表处理或确认。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"}}'

检查队列计数仍为二和一,并且 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

结束环境之前,运行清理检查。

总结

你将一个 SNS 主题连接到两个独立 SQS 队列,将投递来源限定为一个精确主题,并使用消息属性过滤分析通知。实际通知封装证明了共同的发布事件和独立的队列副本。你还区分了 SNS 接受发布与通知到达目标,随后移除了自己的资源,同时保留参考数据。

队列消费者仍需处理重复投递。下一个实验将使用原子 DynamoDB 条件,防止重复任务产生重复业务效果。