连接FIFO Topic的SQS FIFO队列无法接收后续消息
问题背景
使用AWS SQS/SNS FIFO变体做系统事件管理,基于Localstack进行本地测试。每个实体对应一个FIFO Topic,以Project(可部署网站)为例:
- 消息体为JSON格式
MessageGroupId设为实体ID,保证单实体消息顺序处理MessageDeduplicationId用事件类型+实体ID生成,避免重复消息
测试时给同一个Project FIFO Topic订阅了标准队列和FIFO队列,发现:
- 标准队列能接收所有重复请求的消息
- FIFO队列仅收到第一条消息,后续消息无法接收,导致同一项目的重复请求无法处理
附相关代码与日志
Go代码:生成消息分组与去重ID
func generateSNSFIFOMessageIDs(entityID string, eventType string) (string, string) { messageGroupId := entityID // 实体ID作为分组ID,保证单实体消息有序 messageDeduplicationId := fmt.Sprintf("%s-%s", eventType, entityID) // 事件类型+实体ID作为去重ID return messageGroupId, messageDeduplicationId }
CDK代码:创建SNS/SQS资源
import * as sns from 'aws-cdk-lib/aws-sns'; import * as sqs from 'aws-cdk-lib/aws-sqs'; import * as subscriptions from 'aws-cdk-lib/aws-sns-subscriptions'; import { Stack, StackProps } from 'aws-cdk-lib'; import { Construct } from 'constructs'; export class EventStack extends Stack { constructor(scope: Construct, id: string, props?: StackProps) { super(scope, id, props); // 创建Project FIFO Topic const projectFifoTopic = new sns.Topic(this, 'ProjectFifoTopic', { topicName: 'project-events.fifo', fifo: true, contentBasedDeduplication: false, // 手动指定去重ID }); // 标准队列 const standardQueue = new sqs.Queue(this, 'ProjectStandardQueue', { queueName: 'project-standard-queue', }); // FIFO队列 const fifoQueue = new sqs.Queue(this, 'ProjectFifoQueue', { queueName: 'project-events-queue.fifo', fifo: true, contentBasedDeduplication: false, }); // 订阅到标准队列 projectFifoTopic.addSubscription(new subscriptions.SqsSubscription(standardQueue)); // 订阅到FIFO队列(开启原始消息投递) projectFifoTopic.addSubscription(new subscriptions.SqsSubscription(fifoQueue, { rawMessageDelivery: true, })); } }
测试日志
# 标准队列接收(可重复获取到消息) aws --endpoint-url=http://localhost:4566 sqs receive-message --queue-url http://localhost:4566/000000000000/project-standard-queue # 返回第一条消息 aws --endpoint-url=http://localhost:4566 sqs receive-message --queue-url http://localhost:4566/000000000000/project-standard-queue # 返回第二条消息 # FIFO队列接收(仅第一条有返回,后续无) aws --endpoint-url=http://localhost:4566 sqs receive-message --queue-url http://localhost:4566/000000000000/project-events-queue.fifo # 返回第一条消息,包含Attributes: {"MessageGroupId": "proj-123", "MessageDeduplicationId": "DEPLOY-proj-123"} aws --endpoint-url=http://localhost:4566 sqs receive-message --queue-url http://localhost:4566/000000000000/project-events-queue.fifo # 返回空{}
排查方向与解决方案
1. 检查FIFO队列的消息确认状态
SQS FIFO队列中,同一MessageGroupId的消息是顺序处理的——只有前一条消息被确认删除后,下一条消息才会被投递。如果用receive-message获取消息后没有执行delete-message,这条消息会留在队列中(处于隐藏状态),后续同分组消息会被阻塞。
解决步骤:
- 用
delete-message命令删除已接收的消息:
aws --endpoint-url=http://localhost:4566 sqs delete-message --queue-url http://localhost:4566/000000000000/project-events-queue.fifo --receipt-handle "[你的ReceiptHandle值]"
- 删除后再次尝试接收消息,看是否能获取到后续消息。
2. 调整MessageDeduplicationId生成逻辑
当前的去重ID是事件类型+实体ID,同一项目的同一事件类型重复请求会生成完全相同的去重ID。SNS FIFO Topic会过滤掉重复的去重ID消息,不会推送到订阅的FIFO队列;而Localstack对标准队列订阅可能跳过了去重检查,所以标准队列能收到所有消息。
如果重复请求需要被正常处理(不属于重复消息),需要修改去重ID生成规则,加入唯一标识(比如请求ID、时间戳):
func generateSNSFIFOMessageIDs(entityID string, eventType string, requestID string) (string, string) { messageGroupId := entityID // 加入请求ID保证去重ID唯一 messageDeduplicationId := fmt.Sprintf("%s-%s-%s", eventType, entityID, requestID) return messageGroupId, messageDeduplicationId }
3. 验证Localstack版本与实现差异
Localstack对AWS服务的实现可能存在版本差异,建议:
- 升级到最新版本的Localstack,重新测试
- 在原生AWS环境中复现相同配置,确认是否存在同样问题。如果原生AWS正常,说明是Localstack的Bug,可提交Issue反馈。
4. 检查CDK订阅配置
确保FIFO队列的订阅开启了rawMessageDelivery(已在CDK代码中配置),避免SNS包装消息导致的潜在解析问题;同时确认SNS Topic和SQS队列的Region与Localstack一致(默认us-east-1)。
总结
优先排查消息确认是否完成,这是FIFO队列顺序处理的核心规则;其次检查去重ID的生成逻辑是否符合业务需求;最后验证Localstack的实现差异。
内容的提问来源于stack exchange,提问作者Nikola-Milovic

