【问题标题】:SQS FIFO Using MessageGroupId to receive messageSQS FIFO 使用 MessageGroupId 接收消息
【发布时间】:2017-10-03 17:53:08
【问题描述】:

如何使用 messagegroupid 参数仅接收带有我需要的 id 标记的队列消息?

我一直在尝试使用下面的行来检索,但它总是会收到来自其他组 id 的所有队列消息。

List<Message> messages = sqs.receiveMessage(receiveMessageRequest.withAttributeNames("MessageGroupId")).getMessages();

正确的做法应该是什么?

【问题讨论】:

    标签: java amazon-web-services amazon-sqs


    【解决方案1】:

    ReceiveMessageRequest 不用于基于消息属性的过滤。如果您查看ReceiveMessageRequest.html.withAttributeNames() 的文档,它会说:

    需要与每条消息一起返回的属性列表。

    一般来说,您无法过滤从 SQS 返回的消息。您可以限制数量,但不能说,例如,“给我所有符合此模式的消息”。

    【讨论】:

    • 感谢您的回复!这是否意味着我应该有 2 个单独的队列来处理具有不同 messagegroupid 的消息?这样我就不会收到用于其他组 ID 的消息。
    • @JustStarted - 这将是解决问题的一个非常简单的方法。一般来说,创建队列很容易,并且能够对数据进行分区可以简化设计。
    【解决方案2】:

    我的解决方案是利用 ChangeMes​​sageVisibilityBatchRequest (see docs),它本质上将消息发送回队列以进行重新处理。

    我的 lambda 有一个基于时间的触发器。每次它打开时,我都会收到消息,直到不再有消息为止。对于每个批次,我按 MessageGroupId 对消息进行分组,处理并删除第一个组,并将剩余的组消息发送回队列以在下一次迭代中提取。

    这是我的代码的要点:

    (注意_awsSqsService.SendBackToQueueAsync(groupMessages) 方法最终通过 AWS SQS ChangeMessageVisibilityBatchRequest 将消息发送回队列)

    public async Task Run()
    {
        var batchContainsMessages = true;
        while (batchContainsMessages)
        {
            var messageBatch = await _awsSqsService.GetMessageBatchAsync();
            if(messageBatch.Messages.Count > 0 && messageBatch.HttpStatusCode == HttpStatusCode.OK)
            {
                await ProcessMessageBatchAsync(messageBatch.Messages);
            }
            else
            {
                batchContainsMessages = false;
            }
        }
    }
    
    private async Task ProcessMessageBatchAsync(List<Message> messages)
    {
        // SQS fifo queues will often return a batch of messages with different MessageGroupIds.
        // Because of this, we need to group them ourselves, process one group (the first group), 
        // and send the rest back to the queue to be processed in the next iteration. 
        // This ensures that we process as many messages as possible per group in a single batch (max is 10)
        var messageGroups = GetMessageGroups(messages);
    
        var isFirstGroup = true;
        foreach (var group in messageGroups)
        {
            var groupId = Int32.Parse(group.Key);
            var groupMessages = group.Value;
            if (isFirstGroup)
            {
                isFirstGroup = false;
                await ProcessMessagesAsync(groupId, groupMessages);
                await _awsSqsService.DeleteMessagesAsync(groupMessages);
            }
            else
            {
                await _awsSqsService.SendBackToQueueAsync(groupMessages);
            }
        }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-12-31
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-04-09
      • 1970-01-01
      • 2020-10-23
      • 2023-03-03
      相关资源
      最近更新 更多