Python操作AWS SQS FIFO:已删除消息仍重复接收问题排查
AWS SQS FIFO队列消息删除后重复出现,delete_message返回200仍有消息重新入队
我创建了AWS SQS FIFO队列,发送10条消息后执行接收并删除操作,但消息会重复出现。首次接收10条,之后依次收到5条、2条,直到0条;delete_message接口返回200状态码。已设置非零VisibilityTimeout,但问题依旧。相关代码和队列配置如下:
原代码
import json import os import uuid import boto3 from dotenv import load_dotenv load_dotenv() sqs = boto3.client("sqs", os.getenv("AWS_REGION")) receive_queue = os.getenv("RECEIVE_QUEUE_URL") def receive(attempt_id, max_num_messages): response = sqs.receive_message( QueueUrl=receive_queue, ReceiveRequestAttemptId=attempt_id, MaxNumberOfMessages=max_num_messages, VisibilityTimeout=100, WaitTimeSeconds=20, ) if "Messages" not in response: return None, True messages = [message["Body"] for message in response["Messages"]] receipt_handles = [message["ReceiptHandle"] for message in response["Messages"]] print(f"{len(messages)} msgs received") for receipt_handle in receipt_handles: sqs.delete_message(QueueUrl=receive_queue, ReceiptHandle=receipt_handle) receipt_handles.remove(receipt_handle) print(f"{len(receipt_handles)} msgs deleted") return messages, False def send_to_queue(queue_url, data, message_group_id): response = sqs.send_message( QueueUrl=queue_url, MessageBody=json.dumps(data), MessageGroupId=message_group_id, ) return response # 10 msgs created for i in range(10): sent_response = send_to_queue( receive_queue, {"key": i}, message_group_id=os.getenv("MESSAGE_GROUP_ID") ) # 10 msgs received receive(str(uuid.uuid4()), 10) # Now still 5 msgs are "inflight" and received again
队列配置

问题根源
问题出在receive函数的循环逻辑:
for receipt_handle in receipt_handles: sqs.delete_message(QueueUrl=receive_queue, ReceiptHandle=receipt_handle) receipt_handles.remove(receipt_handle)
遍历列表的同时调用remove()会导致迭代器跳过元素。比如初始列表有10个句柄,第一次循环处理第1个并删除后,列表长度变为9;下一次迭代会直接取原列表的第3个元素(迭代指针已移动),跳过第2个句柄。最终仅一半消息被实际删除,剩余消息在VisibilityTimeout到期后重新回到队列,导致重复接收。
另外,print(f"{len(receipt_handles)} msgs deleted")是错误的——这里打印的是剩余未删除的句柄数,而非已删除数量,会误导你以为所有消息都被处理。
修复方案
1. 修改循环逻辑,避免遍历同时修改原列表
可以遍历列表的副本,或者直接计数已删除的消息:
def receive(attempt_id, max_num_messages): response = sqs.receive_message( QueueUrl=receive_queue, ReceiveRequestAttemptId=attempt_id, MaxNumberOfMessages=max_num_messages, VisibilityTimeout=100, WaitTimeSeconds=20, ) if "Messages" not in response: return None, True messages = [message["Body"] for message in response["Messages"]] receipt_handles = [message["ReceiptHandle"] for message in response["Messages"]] print(f"{len(messages)} msgs received") deleted_count = 0 # 遍历原列表,无需修改它 for receipt_handle in receipt_handles: sqs.delete_message(QueueUrl=receive_queue, ReceiptHandle=receipt_handle) deleted_count += 1 print(f"{deleted_count} msgs deleted") return messages, False
2. 验证逻辑
修改后,所有接收的消息句柄都会被遍历并执行删除操作,deleted_count会等于接收的消息数量,确认所有消息都被彻底删除,不会在VisibilityTimeout后重新入队。
内容的提问来源于stack exchange,提问作者JTX
相关产品推荐
相关产品推荐

