【发布时间】:2021-01-22 08:39:22
【问题描述】:
我在 SQS 中遇到了一个奇怪的问题。让我简化我的用例,我的 FIFO 队列中有 7 条消息,我的独立应用程序应该不断地以相同的顺序轮询消息,以用于我的业务案例。比如我的app读取了message1,经过一些业务处理后,app会删除它并将同一条消息重新发送到同一个队列(队列的尾部),这些步骤将无休止地继续下一组消息。在这里,我的期望是我的应用程序将连续轮询消息并根据队列中的消息以相同的顺序执行操作,但这就是问题出现的地方。第一次从队列中读取消息时,将其删除,然后将同一条消息重新发送到同一个队列中,即使在成功发送消息结果之后,重新发送的消息也不存在于队列中。
我已经包含下面的代码来模拟问题,简单地说,Test_Queue.fifo 队列与Test_Queue_DLQ.fifo 配置为reDrivePolicy 创建。在创建队列后的第一次,消息被发布 -> "Test_Message" 到 Test_Queue.fifo 队列(在响应中获取 MessageId )并长轮询队列以读取消息,并在迭代 ReceiveMessageResult#getMessages 之后,删除消息(获取 MessageId 作为响应)。同样,在成功删除消息后,将相同的消息重新发送到同一队列的尾部(获取 MessageId 作为响应)。但是,重新发布的消息不在队列中。当我检查 AWS 管理控制台时,Messages available 和 Messages in flight 部分中的消息计数为 0,而Test_Queue_DLQ.fifo queue 中甚至不存在重新发布的消息。根据 SQS 文档,如果我们删除该消息,即使它在飞行模式下也应该被删除,因此重新发布相同的消息应该不是问题。我怀疑在 SQS 方面,他们正在执行一些相等比较并在 visibleTimeOut 间隔期间跳过相同的消息,以避免在分布式环境中对相同的消息进行重复数据删除,但无法获得任何清晰的画面。
代码 sn-p 模拟上述问题
public class SQSIssue {
@Test
void sqsMessageAbsenceIssueTest() {
AmazonSQS amazonSQS = AmazonSQSClientBuilder.standard().withEndpointConfiguration(new AwsClientBuilder
.EndpointConfiguration("https://sqs.us-east-2.amazonaws.com", "us-east-2"))
.withCredentials(new AWSStaticCredentialsProvider(new BasicAWSCredentials(
"<accessKey>", "<secretKey>"))).build();
//create queue
String queueUrl = createQueues(amazonSQS);
String message = "Test_Message";
String groupId = "Group1";
//Sending message -> "Test_Message"
sendMessage(amazonSQS, queueUrl, message, groupId);
//Reading the message and deleting using message.getReceiptHandle()
readAndDeleteMessage(amazonSQS, queueUrl);
//Reposting the same message into the queue -> "Test_Message"
sendMessage(amazonSQS, queueUrl, message, groupId);
ReceiveMessageRequest receiveMessageRequest = new ReceiveMessageRequest()
.withQueueUrl(queueUrl)
.withWaitTimeSeconds(5)
.withMessageAttributeNames("All")
.withVisibilityTimeout(30)
.withMaxNumberOfMessages(10);
ReceiveMessageResult receiveMessageResult = amazonSQS.receiveMessage(receiveMessageRequest);
//Here I am expecting the message presence in the queue as I recently reposted the same message into the same queue after the message deletion
Assertions.assertFalse(receiveMessageResult.getMessages().isEmpty());
}
private void readAndDeleteMessage(AmazonSQS amazonSQS, String queueUrl) {
ReceiveMessageRequest receiveMessageRequest = new ReceiveMessageRequest()
.withQueueUrl(queueUrl)
.withWaitTimeSeconds(5)
.withMessageAttributeNames("All")
.withVisibilityTimeout(30)
.withMaxNumberOfMessages(10);
ReceiveMessageResult receiveMessageResult = amazonSQS.receiveMessage(receiveMessageRequest);
receiveMessageResult.getMessages().forEach(message -> amazonSQS.deleteMessage(queueUrl, message.getReceiptHandle()));
}
private String createQueues(AmazonSQS amazonSQS) {
String queueName = "Test_Queue.fifo";
String deadLetterQueueName = "Test_Queue_DLQ.fifo";
//Creating DeadLetterQueue
CreateQueueRequest createDeadLetterQueueRequest = new CreateQueueRequest()
.addAttributesEntry("FifoQueue", "true")
.addAttributesEntry("ContentBasedDeduplication", "true")
.addAttributesEntry("VisibilityTimeout", "600")
.addAttributesEntry("MessageRetentionPeriod", "262144");
createDeadLetterQueueRequest.withQueueName(deadLetterQueueName);
CreateQueueResult createDeadLetterQueueResult = amazonSQS.createQueue(createDeadLetterQueueRequest);
GetQueueAttributesResult getQueueAttributesResult = amazonSQS.getQueueAttributes(
new GetQueueAttributesRequest(createDeadLetterQueueResult.getQueueUrl())
.withAttributeNames("QueueArn"));
String deadLetterQueueArn = getQueueAttributesResult.getAttributes().get("QueueArn");
//Creating Actual Queue with DeadLetterQueue configured
CreateQueueRequest createQueueRequest = new CreateQueueRequest()
.addAttributesEntry("FifoQueue", "true")
.addAttributesEntry("ContentBasedDeduplication", "true")
.addAttributesEntry("VisibilityTimeout", "600")
.addAttributesEntry("MessageRetentionPeriod", "262144");
createQueueRequest.withQueueName(queueName);
String reDrivePolicy = "{\"maxReceiveCount\":\"5\", \"deadLetterTargetArn\":\""
+ deadLetterQueueArn + "\"}";
createQueueRequest.addAttributesEntry("RedrivePolicy", reDrivePolicy);
CreateQueueResult createQueueResult = amazonSQS.createQueue(createQueueRequest);
return createQueueResult.getQueueUrl();
}
private void sendMessage(AmazonSQS amazonSQS, String queueUrl, String message, String groupId) {
SendMessageRequest sendMessageRequest = new SendMessageRequest()
.withQueueUrl(queueUrl)
.withMessageBody(message)
.withMessageGroupId(groupId);
SendMessageResult sendMessageResult = amazonSQS.sendMessage(sendMessageRequest);
Assertions.assertNotNull(sendMessageResult.getMessageId());
}
}
【问题讨论】:
标签: amazon-web-services aws-sdk amazon-sqs