Skip to content

Commit 05daa45

Browse files
authored
Enhance sqs partial acknowledgement handling for listeners (#1562)
* Fix : Add logging for SQS Partial Acknowledgement * Style : Modify code style to suit the characteristics of the existing repository * Fix : Test code changes due to added logging - Changed null to DeleteMessageBatchResponse * Test : Add Test for partialBatchFailure * Fix : Added logging of failed ID list * Fix: Added success/failure lists to SqsAcknowledgementException DeleteMessageBatch response - Correlate items using message IDs for accurate mapping - Pass different lists to the SqsAcknowledgementException constructor to handle partial failures * Refactoring: Extract the createPartialFailureException helper method - Extract partial error handling logic from the deleteMessages method - Improved code readability and maintainability - placed the calling method at the top and the helper method at the bottom * Refactoring: Separate the batch response processing deletion part from SqsAcknowledgementExecutor. * Fix: improve partial SQS acknowledgement failure handling - handle DeleteMessageBatch partial failures in a dedicated method - map AWS failed entry IDs to original message IDs - populate SqsAcknowledgementException with successful/failed message lists - avoid wrapping SqsAcknowledgementException twice - extend SqsAcknowledgementExecutorTests to assert partial failure mapping * Docs: document partial acknowledgement failure handling in callbacks - fix callback interface name to AsyncAcknowledgementResultCallback - explain that onFailure receives SqsAcknowledgementException on partial failures - add an example using successful/failed acknowledgement message lists for retry * Fix: add fail-safe handling for uncorrelated SQS acknowledgement failure ids * Test: add coverage for uncorrelated acknowledgement failure ids * Docs: document fail-safe correlation behavior in acknowledgement result callbacks --------- Co-authored-by: joyoungjae <jaonz6057@gmail.com>
1 parent 3a07c0e commit 05daa45

3 files changed

Lines changed: 152 additions & 7 deletions

File tree

docs/src/main/asciidoc/sqs.adoc

Lines changed: 18 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -468,6 +468,7 @@ SqsTemplate.builder().configure(options -> options.acknowledgementMode(TemplateA
468468
```
469469

470470
If an error occurs during acknowledgement, a `SqsAcknowledgementException` is thrown, containing both the messages that were successfully acknowledged and those which failed.
471+
See <<sqs-acknowledgement-result-callback>> for details on inspecting partial failure results and the fail-safe correlation behavior.
471472

472473
To acknowledge messages received with `MANUAL` acknowledgement, the `Acknowledgement#acknowledge` and `Acknowledgement#acknowledgeAsync` methods can be used.
473474

@@ -2018,9 +2019,10 @@ NOTE: PARALLEL is the default for FIFO because ordering is guaranteed for proces
20182019
This assures no messages from a given `MessageGroup` will be polled until the previous batch is acknowledged.
20192020
Implementations of this interface will be executed after an acknowledgement execution completes with either success or failure.
20202021

2022+
[[sqs-acknowledgement-result-callback]]
20212023
==== Acknowledgement Result Callback
20222024

2023-
The framework offers the `AcknowledgementResultCallback` and `AsyncAcknowledgementCallback` interfaces that can be added to a `SqsMessageListenerContainer` or `SqsMessageListenerContainerFactory`.
2025+
The framework offers the `AcknowledgementResultCallback` and `AsyncAcknowledgementResultCallback` interfaces that can be added to a `SqsMessageListenerContainer` or `SqsMessageListenerContainerFactory`.
20242026

20252027
```java
20262028
public interface AcknowledgementResultCallback<T> {
@@ -2048,6 +2050,21 @@ public interface AsyncAcknowledgementResultCallback<T> {
20482050
}
20492051
```
20502052

2053+
If an acknowledgement operation partially fails, for example when `DeleteMessageBatch` returns failed entries, the callback `onFailure` receives a `SqsAcknowledgementException`.
2054+
Use `getSuccessfullyAcknowledgedMessages()` and `getFailedAcknowledgementMessages()` to inspect the acknowledgement result and retry only failed messages if needed.
2055+
2056+
```java
2057+
@Override
2058+
public void onFailure(Collection<Message<Object>> messages, Throwable t) {
2059+
if (t instanceof SqsAcknowledgementException ex) {
2060+
Collection<Message<?>> failedMessages = ex.getFailedAcknowledgementMessages();
2061+
// retry only failedMessages
2062+
}
2063+
}
2064+
```
2065+
2066+
NOTE: If the failure IDs returned by AWS cannot be correlated with the original request IDs, a fail-safe is applied: `getSuccessfullyAcknowledgedMessages()` returns an empty collection and `getFailedAcknowledgementMessages()` returns all messages in the batch to prevent silent misclassification.
2067+
20512068
```java
20522069
@Bean
20532070
public SqsMessageListenerContainerFactory<Object> defaultSqsListenerContainerFactory(SqsAsyncClient sqsAsyncClient) {

spring-cloud-aws-sqs/src/main/java/io/awspring/cloud/sqs/listener/acknowledgement/SqsAcknowledgementExecutor.java

Lines changed: 58 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -22,9 +22,11 @@
2222
import io.awspring.cloud.sqs.listener.QueueAttributesAware;
2323
import io.awspring.cloud.sqs.listener.SqsAsyncClientAware;
2424
import io.awspring.cloud.sqs.listener.SqsHeaders;
25+
import java.util.ArrayList;
2526
import java.util.Collection;
2627
import java.util.Collections;
27-
import java.util.UUID;
28+
import java.util.List;
29+
import java.util.Set;
2830
import java.util.concurrent.CompletableFuture;
2931
import java.util.concurrent.CompletionException;
3032
import java.util.stream.Collectors;
@@ -34,8 +36,10 @@
3436
import org.springframework.util.Assert;
3537
import org.springframework.util.StopWatch;
3638
import software.amazon.awssdk.services.sqs.SqsAsyncClient;
39+
import software.amazon.awssdk.services.sqs.model.BatchResultErrorEntry;
3740
import software.amazon.awssdk.services.sqs.model.DeleteMessageBatchRequest;
3841
import software.amazon.awssdk.services.sqs.model.DeleteMessageBatchRequestEntry;
42+
import software.amazon.awssdk.services.sqs.model.DeleteMessageBatchResponse;
3943

4044
/**
4145
* {@link AcknowledgementExecutor} implementation for SQS queues. Handle the messages deletion, usually requested by an
@@ -95,12 +99,61 @@ private CompletableFuture<Void> deleteMessages(Collection<Message<T>> messagesTo
9599
StopWatch watch = new StopWatch();
96100
watch.start();
97101
return CompletableFutures.exceptionallyCompose(this.sqsAsyncClient
98-
.deleteMessageBatch(createDeleteMessageBatchRequest(messagesToAck))
99-
.thenRun(() -> {}),
100-
t -> CompletableFutures.failedFuture(createAcknowledgementException(messagesToAck, t)))
102+
.deleteMessageBatch(createDeleteMessageBatchRequest(messagesToAck)).thenCompose(
103+
response -> handleDeleteMessageBatchResponse(messagesToAck, response)),
104+
t -> toAcknowledgementFailure(messagesToAck, t))
101105
.whenComplete((v, t) -> logAckResult(messagesToAck, t, watch));
102106
}
103107

108+
private CompletableFuture<Void> handleDeleteMessageBatchResponse(Collection<Message<T>> messagesToAck,
109+
DeleteMessageBatchResponse response) {
110+
if (!response.failed().isEmpty()) {
111+
return CompletableFutures.<Void>failedFuture(createPartialFailureException(messagesToAck, response));
112+
}
113+
return CompletableFuture.<Void>completedFuture(null);
114+
}
115+
116+
private CompletableFuture<Void> toAcknowledgementFailure(Collection<Message<T>> messagesToAck, Throwable throwable) {
117+
Throwable cause = throwable instanceof CompletionException && throwable.getCause() != null ? throwable.getCause()
118+
: throwable;
119+
if (cause instanceof SqsAcknowledgementException) {
120+
return CompletableFutures.<Void>failedFuture(cause);
121+
}
122+
return CompletableFutures.<Void>failedFuture(createAcknowledgementException(messagesToAck, cause));
123+
}
124+
125+
private SqsAcknowledgementException createPartialFailureException(Collection<Message<T>> messages,
126+
DeleteMessageBatchResponse response) {
127+
Set<String> messageIds = messages.stream().map(MessageHeaderUtils::getId).collect(Collectors.toSet());
128+
Set<String> failedIds = response.failed().stream()
129+
.map(BatchResultErrorEntry::id)
130+
.collect(Collectors.toSet());
131+
132+
if (!messageIds.containsAll(failedIds)) {
133+
logger.warn("Could not correlate all acknowledgement failure ids in queue {}: {}", this.queueName,
134+
failedIds);
135+
return new SqsAcknowledgementException("Could not correlate acknowledgement failure ids: " + failedIds,
136+
Collections.emptyList(), messages.stream().map(msg -> (Message<?>) msg).collect(Collectors.toList()),
137+
this.queueUrl, null);
138+
}
139+
140+
List<Message<?>> successfulMessages = new ArrayList<>();
141+
List<Message<?>> failedMessages = new ArrayList<>();
142+
143+
for(Message<T> msg : messages) {
144+
if(failedIds.contains(MessageHeaderUtils.getId(msg))) {
145+
failedMessages.add(msg);
146+
} else {
147+
successfulMessages.add(msg);
148+
}
149+
}
150+
151+
logger.warn("Some messages could not be acknowledged in queue {}: {}", this.queueName, failedIds);
152+
153+
return new SqsAcknowledgementException("Error acknowledging messages " + failedIds, successfulMessages,
154+
failedMessages, this.queueUrl, null);
155+
}
156+
104157
private DeleteMessageBatchRequest createDeleteMessageBatchRequest(Collection<Message<T>> messagesToAck) {
105158
return DeleteMessageBatchRequest
106159
.builder()
@@ -113,7 +166,7 @@ private DeleteMessageBatchRequestEntry toDeleteMessageEntry(Message<T> message)
113166
return DeleteMessageBatchRequestEntry
114167
.builder()
115168
.receiptHandle(MessageHeaderUtils.getHeaderAsString(message, SqsHeaders.SQS_RECEIPT_HANDLE_HEADER))
116-
.id(UUID.randomUUID().toString())
169+
.id(MessageHeaderUtils.getId(message))
117170
.build();
118171
}
119172
// @formatter:on

spring-cloud-aws-sqs/src/test/java/io/awspring/cloud/sqs/listener/acknowledgement/SqsAcknowledgementExecutorTests.java

Lines changed: 76 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
import io.awspring.cloud.sqs.listener.SqsHeaders;
2929
import java.util.Collection;
3030
import java.util.Collections;
31+
import java.util.List;
3132
import java.util.concurrent.CompletableFuture;
3233
import java.util.concurrent.CompletionException;
3334
import org.junit.jupiter.api.Test;
@@ -38,8 +39,10 @@
3839
import org.springframework.messaging.Message;
3940
import org.springframework.messaging.MessageHeaders;
4041
import software.amazon.awssdk.services.sqs.SqsAsyncClient;
42+
import software.amazon.awssdk.services.sqs.model.BatchResultErrorEntry;
4143
import software.amazon.awssdk.services.sqs.model.DeleteMessageBatchRequest;
4244
import software.amazon.awssdk.services.sqs.model.DeleteMessageBatchRequestEntry;
45+
import software.amazon.awssdk.services.sqs.model.DeleteMessageBatchResponse;
4346

4447
/**
4548
* Tests for {@link SqsAcknowledgementExecutor}.
@@ -58,23 +61,31 @@ class SqsAcknowledgementExecutorTests {
5861
@Mock
5962
Message<String> message;
6063

64+
@Mock
65+
Message<String> secondMessage;
66+
6167
String queueName = "sqsAcknowledgementExecutorTestsQueueName";
6268

6369
String queueUrl = "sqsAcknowledgementExecutorTestsQueueUrl";
6470

6571
String receiptHandle = "sqsAcknowledgementExecutorTestsQueueReceiptHandle";
6672

73+
String secondReceiptHandle = "sqsAcknowledgementExecutorTestsQueueSecondReceiptHandle";
74+
6775
MessageHeaders messageHeaders = new MessageHeaders(
6876
Collections.singletonMap(SqsHeaders.SQS_RECEIPT_HANDLE_HEADER, receiptHandle));
6977

78+
MessageHeaders secondMessageHeaders = new MessageHeaders(
79+
Collections.singletonMap(SqsHeaders.SQS_RECEIPT_HANDLE_HEADER, secondReceiptHandle));
80+
7081
@Test
7182
void shouldDeleteMessages() throws Exception {
7283
Collection<Message<String>> messages = Collections.singletonList(message);
7384
given(message.getHeaders()).willReturn(messageHeaders);
7485
given(queueAttributes.getQueueName()).willReturn(queueName);
7586
given(queueAttributes.getQueueUrl()).willReturn(queueUrl);
7687
given(sqsAsyncClient.deleteMessageBatch(any(DeleteMessageBatchRequest.class)))
77-
.willReturn(CompletableFuture.completedFuture(null));
88+
.willReturn(CompletableFuture.completedFuture(DeleteMessageBatchResponse.builder().build()));
7889

7990
SqsAcknowledgementExecutor<String> executor = new SqsAcknowledgementExecutor<>();
8091
executor.setSqsAsyncClient(sqsAsyncClient);
@@ -127,4 +138,68 @@ void shouldWrapIfErrorIsThrown() {
127138
.extracting(SqsAcknowledgementException::getQueue).isEqualTo(queueUrl);
128139
}
129140

141+
@Test
142+
void shouldWrapPartialBatchFailure() {
143+
Message<String> failedMessage = message;
144+
Message<String> successfulMessage = secondMessage;
145+
MessageHeaders failedMessageHeaders = messageHeaders;
146+
MessageHeaders successfulMessageHeaders = secondMessageHeaders;
147+
Collection<Message<String>> messagesToAck = List.of(failedMessage, successfulMessage);
148+
149+
given(failedMessage.getHeaders()).willReturn(failedMessageHeaders);
150+
given(successfulMessage.getHeaders()).willReturn(successfulMessageHeaders);
151+
given(queueAttributes.getQueueName()).willReturn(queueName);
152+
given(queueAttributes.getQueueUrl()).willReturn(queueUrl);
153+
154+
BatchResultErrorEntry failedEntry = BatchResultErrorEntry.builder().id(failedMessageHeaders.getId().toString())
155+
.code("ReceiptHandleIsInvalid").message("Receipt handle expired").build();
156+
157+
DeleteMessageBatchResponse partialFailureResponse = DeleteMessageBatchResponse.builder().failed(failedEntry)
158+
.build();
159+
160+
given(sqsAsyncClient.deleteMessageBatch(any(DeleteMessageBatchRequest.class)))
161+
.willReturn(CompletableFuture.completedFuture(partialFailureResponse));
162+
163+
SqsAcknowledgementExecutor<String> executor = new SqsAcknowledgementExecutor<>();
164+
executor.setSqsAsyncClient(sqsAsyncClient);
165+
executor.setQueueAttributes(queueAttributes);
166+
167+
assertThatThrownBy(() -> executor.execute(messagesToAck).join()).isInstanceOf(CompletionException.class)
168+
.getCause().isInstanceOf(SqsAcknowledgementException.class)
169+
.asInstanceOf(type(SqsAcknowledgementException.class)).satisfies(ex -> {
170+
assertThat(ex.getFailedAcknowledgementMessages()).containsExactly(failedMessage);
171+
assertThat(ex.getSuccessfullyAcknowledgedMessages()).containsExactly(successfulMessage);
172+
});
173+
}
174+
175+
@Test
176+
void shouldTreatAllMessagesAsFailedIfAwsFailureIdCannotBeCorrelated() {
177+
Collection<Message<String>> messagesToAck = List.of(message, secondMessage);
178+
179+
given(message.getHeaders()).willReturn(messageHeaders);
180+
given(secondMessage.getHeaders()).willReturn(secondMessageHeaders);
181+
given(queueAttributes.getQueueName()).willReturn(queueName);
182+
given(queueAttributes.getQueueUrl()).willReturn(queueUrl);
183+
184+
BatchResultErrorEntry failedEntry = BatchResultErrorEntry.builder().id("unknown-id")
185+
.code("ReceiptHandleIsInvalid").message("Receipt handle expired").build();
186+
187+
DeleteMessageBatchResponse partialFailureResponse = DeleteMessageBatchResponse.builder().failed(failedEntry)
188+
.build();
189+
190+
given(sqsAsyncClient.deleteMessageBatch(any(DeleteMessageBatchRequest.class)))
191+
.willReturn(CompletableFuture.completedFuture(partialFailureResponse));
192+
193+
SqsAcknowledgementExecutor<String> executor = new SqsAcknowledgementExecutor<>();
194+
executor.setSqsAsyncClient(sqsAsyncClient);
195+
executor.setQueueAttributes(queueAttributes);
196+
197+
assertThatThrownBy(() -> executor.execute(messagesToAck).join()).isInstanceOf(CompletionException.class)
198+
.getCause().isInstanceOf(SqsAcknowledgementException.class)
199+
.asInstanceOf(type(SqsAcknowledgementException.class)).satisfies(ex -> {
200+
assertThat(ex.getSuccessfullyAcknowledgedMessages()).isEmpty();
201+
assertThat(ex.getFailedAcknowledgementMessages()).containsExactlyInAnyOrder(message, secondMessage);
202+
});
203+
}
204+
130205
}

0 commit comments

Comments
 (0)