使用 SQS 发送和消费任务

AWSBeginner
立即练习

介绍

一个订单服务需要现在接受任务,稍后再处理。你将创建队列、领取任务、运行提供的工作程序,并在确认消息之前检查存储的订单。

请先完成 通过 Lambda 读取和写入 DynamoDB 及其前置引导实验。这个独立 VM 提供工作程序、限定范围的角色、空订单表和参考数据;没有预先准备队列或任务。

创建独立的任务队列

本步骤中,在发送任何任务之前,创建空队列并检查提供的工作程序。

使用 Terminal 旁的 AWS View 对比当前队列、工作程序结果和存储的订单。保留无关的参考数据。

Amazon Simple Queue Service(SQS) 将任务存储为消息,供消费者处理。生产者发送任务,消费者处理任务。队列将两者的时间分离:生产者不必等待消费者完成。Standard 队列可能多次投递同一条消息;领取和确认是独立操作。

从准备好的工作目录开始:

cd /home/labex/project

检查提供的工作程序和空订单表。工作程序按每件商品 250 美分加 100 美分手续费计算金额。其执行角色只能写入订单表;参考表是需要保留的无关数据。

aws lambda get-function-configuration --function-name labex-q01-worker --query '{Name:FunctionName,Role:Role,Runtime:Runtime}'
aws dynamodb scan --table-name labex-q01-orders --query Items

预期项目列表为空。创建自己的队列。--query QueueUrl --output text 将队列地址选为纯文本;$(...) 将该输出存入 shell 变量 QUEUE_URL,供后续命令使用。

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

可见性超时让消费者有 300 秒处理时间,之后已领取的消息可能再次变得可用。这是临时隐藏,不是删除。下一个实验将探索超时到期与再次投递。

检查队列身份和当前计数:

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

预期 VisibilityTimeout 为 300,两个消息计数均为 0。计数是近似的运行信号,不能保证业务已完成。在 AWS View 中,你的队列应显示零个可用任务和零个处理中的任务;提供的 Lambda 没有执行记录,订单表仍为空。

Amazon SQS 官方 Console 队列详情示例

官方 Console 展示的队列名称、类型、URL 和 ARN,与 CLI 检查的字段相同。这些是示例值;请继续使用自己的 QUEUE_URL。使用 AWS View 观察本实验的资源。

来源:AWS SQS。

发送并领取订单任务

本步骤中,发送并领取任务,观察可用消息与处理中的消息之间的区别。

队列投递与业务结果

先领取,再处理,最后确认:确认订单已存储后,使用当前 receipt handle 删除队列消息。

消息正文(body)是应用数据。SQS 将 JSON 保存为文本;消费者必须解释它。发送一个小型模拟订单任务:

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

响应包含 MessageId 和 MD5OfMessageBody。消息 ID 标识消息;MD5 是正文字节的摘要。这表示队列已接受消息,不表示订单已完成。AWS View 此时显示一个可用任务,而订单表仍为空。

领取一条消息并保存响应。> 将输出重定向到 received.json,而不是显示出来。--wait-time-seconds 5 允许在没有立即可用任务时进行短暂的长轮询等待。

aws sqs receive-message --queue-url "$QUEUE_URL" --max-number-of-messages 1 --wait-time-seconds 5 --message-system-attribute-names ApproximateReceiveCount --output json > received.json

使用 jq 读取安全的应用正文和领取次数;它用于从 JSON 中选择字段:

jq '.Messages[0] | {MessageId,Body,Attributes}' received.json

预期正文包含 queue-order 和数量 2,第一次领取时的领取次数为 1。完整响应还包含 receipt handle(接收句柄):它属于这一次特定的投递尝试。确认此消息时,你将使用最新的 handle。

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

预期可用消息为 0,不可见消息为 1。在 AWS View 中,任务处于处理中(in flight),但仍然没有订单。领取没有处理或删除它。请在 300 秒内继续下一步。如果阅读时间更长,在处理和删除之前,再次领取并保存到同一文件,以获取最新的 receipt handle。

AWS View 示例:一个任务处于处理中,但尚无工作程序执行记录或存储的订单

这个实际示例展示处理前的投递状态。你工作目录中的消息 ID 会有所不同。

在确认之前处理任务

本步骤中,处理已领取的正文,验证持久化订单,然后确认消息。

使用实际领取的正文作为工作程序输入。fromjson 将 SQS 响应中的 JSON 字符串转换为 JSON 对象;重定向将该对象写入 job.json。

jq '.Messages[0].Body | fromjson' received.json > job.json

提供的 Lambda 接受这个订单对象。与 Lambda 课程一样,fileb://job.json 发送文件字节,末尾的文件名用于接收函数响应。调用工作程序:

aws lambda invoke --function-name labex-q01-worker --payload fileb://job.json worker-response.json

仅有成功的 Invoke API 状态不能证明处理成功。检查响应正文:

cat worker-response.json

预期看到 processed: true、数量 2 和 total_cents: 600。然后读取存储的项目,不要只依赖函数返回的消息:

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

预期看到 queue-order、数量 2 和总金额 600。AWS View 显示实际工作程序输入和结果,以及相同的持久化订单。任务会一直处于处理中,直到你确认它。如果工作程序响应或存储项目有误,应保留消息用于诊断,不要删除它。

确认订单后,将当前 receipt handle 选为纯文本,并从队列中删除这次投递的消息:

RECEIPT_HANDLE=$(jq -r '.Messages[0].ReceiptHandle' received.json)
aws sqs delete-message --queue-url "$QUEUE_URL" --receipt-handle "$RECEIPT_HANDLE"

成功删除通常不会打印任何内容。再次读取计数:

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

两个计数都应为 0。AWS View 显示空队列和一个存储的订单。你已区分队列接受、临时投递、业务处理和确认。在实际应用中,Standard 队列仍可能再次投递消息;这个顺序不构成端到端的恰好一次保证。后续实验将教授持久化的重复保护。

只清理自己的队列和结果

本步骤中,移除你的队列、结果和日志,同时保留提供的资源。

移除队列和你创建的订单,同时保留提供的工作程序、表结构和参考数据。删除队列会丢弃所有剩余任务,因此先确认上一步已成功完成。

aws sqs delete-queue --queue-url "$QUEUE_URL"
aws dynamodb delete-item --table-name labex-q01-orders --key '{"id":{"S":"queue-order"}}'

工作程序执行时创建了 CloudWatch Logs 日志组。将本实验的执行日志一并删除:

aws logs delete-log-group --log-group-name /aws/lambda/labex-q01-worker

通过成功的服务查询检查资源状态:

aws sqs list-queues
aws dynamodb scan --table-name labex-q01-orders --query Items
aws dynamodb scan --table-name labex-q01-reference --query Items

不应再有任何队列 URL,订单项目列表应为空,参考项目必须仍包含 keep unchanged。AWS View 显示没有队列、订单或执行日志,而提供的工作程序和参考资源仍然存在。身份验证或网络错误不能证明删除成功。

检查资源后,移除普通响应文件:

rm -f received.json job.json worker-response.json

结束环境之前,运行本步骤的检查。

总结

你创建了 SQS Standard 队列,发送并领取 JSON 订单任务,验证实际 Lambda 处理与持久化的 DynamoDB 项目,并通过 receipt handle 确认消息。你区分了接受、处理中的投递和业务完成,然后移除了自己的队列、订单和执行日志,同时保留提供的资源。