介绍
同一订单的重复任务必须保留已完成的业务结果。你将为提供的消费者添加原子写入保护,测试重复订单和独立订单,并确认消息已成功确认。
请先完成 使用死信队列隔离失败任务、使用条件写入防止重复订单 和 配置并诊断 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,并创建批大小为一的事件源映射:
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。工作程序确认内容一致的重复消息,而不会再次写入订单。
幂等消费者在重复任务到来时保留相同的业务结果。DynamoDB 在与 PutItem 相同的写入操作中评估 ConditionExpression='attribute_not_exists(id)',避免独立的先读后写竞态。业务 id 是键;如果使用 SQS messageId,生产者新发送的同一订单副本仍会再次写入。
条件失败时,ReturnValuesOnConditionCheckFailure='ALL_OLD' 返回已有项目。代码只捕获 ConditionalCheckFailedException,并将该项目与请求的项目比较。详情一致表示已完成的重复任务;详情冲突或其他错误仍然失败,不会被悄悄确认。
写入下面的完整工作程序。cat > app.py 替换文件,here-document 通过 PY 结束标记提供内容。为 'PY' 加引号可防止 shell 展开 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 副本。
将 app.py 打包到 ZIP 根目录。已有函数的处理程序为 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}'
观察 AWS View,直到队列为空、工作程序报告 processed: true,且订单出现。然后读取它:
aws dynamodb get-item --table-name labex-q05-orders --key '{"id":{"S":"dedup-order"}}' --consistent-read --query Item
预期数量为 4,总金额为 1100。写入决定显示条件、Accepted、没有之前的项目,以及随后存储的项目。仅配置条件或上传文件不能证明任务成功完成。
消费重复任务而不重复写入
本步骤中,发送同一业务订单的两个新副本,并证明另一个订单仍然独立完成。
将相同业务载荷发送两次。每个 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}'
观察 AWS View,直到四次实际执行都已返回且队列为空。两次重复执行报告 processed: false、duplicate: true。DynamoDB 以 ConditionalCheckFailedException 拒绝两次条件写入;每次写入前后的项目都未改变。第一个订单与 other-order 各有一次被接受的写入。消费者确认重复副本,而不是重试已完成的任务。
读取业务结果:
aws dynamodb scan --table-name labex-q05-orders --query Items
预期只有 dedup-order/4/1100 和 other-order/1/350;项目顺序可能不同。检查实际工作程序结果。日志事件可能包含多行;此管道将原生 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
两个计数均为零。实际发生了四次消息投递和函数执行、两次被接受的业务写入,以及两次原生条件拒绝。仅计算表项目数量无法发现无条件覆盖;还要使用实际写入决定及对应的消费者结果。
保护逻辑保护的是一个业务项目。它不会将 DynamoDB 写入和 SQS 确认变成单一原子操作。订单记录是持久化的重复保护。删除它后,同一个键可以再次创建项目,因此保留期限和业务键设计属于实际系统策略的一部分。这些模拟重复任务证明了同一业务任务的安全处理;它们不能证明生产并发或端到端恰好一次保证。

移除消费者连接和业务结果
本步骤中,确认重复任务处理后,移除自己创建的连接和资源。
删除源队列之前,先停止你的事件源映射:
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 业务键条件。第一个订单完成;两个分别发送的副本在原生条件拒绝保留原始项目后被消费并确认。不同的键仍产生了自己的结果。实际写入、工作程序结果和空队列状态,在按限定范围清理之前证明了结果。
挑战将应用有限次数的失败隔离,之后的无服务器项目会将 DLQ 恢复与此业务键保护结合起来。



