Skip to content

Commit 3af9937

Browse files
committed
[Java IO] Stabilize SqsIO timeout batch tests
1 parent 8f4a816 commit 3af9937

1 file changed

Lines changed: 43 additions & 9 deletions

File tree

sdks/java/io/amazon-web-services2/src/test/java/org/apache/beam/sdk/io/aws2/sqs/SqsIOWriteBatchesTest.java

Lines changed: 43 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@
3838

3939
import java.util.Arrays;
4040
import java.util.HashSet;
41+
import java.util.List;
4142
import java.util.Map;
4243
import java.util.Set;
4344
import java.util.stream.Collectors;
@@ -263,10 +264,24 @@ public void testWriteBatchesWithTimeout() {
263264

264265
p.run().waitUntilFinish();
265266

266-
SendMessageBatchRequestEntry[] entries = entries(range(0, 5));
267-
// due to added delay, batches are timed out on arrival of every 3rd msg
268-
verify(sqs).sendMessageBatch(request("queue", entries[0], entries[1], entries[2]));
269-
verify(sqs).sendMessageBatch(request("queue", entries[3], entries[4]));
267+
ArgumentCaptor<SendMessageBatchRequest> captor =
268+
ArgumentCaptor.forClass(SendMessageBatchRequest.class);
269+
verify(sqs, atLeastOnce()).sendMessageBatch(captor.capture());
270+
271+
List<SendMessageBatchRequest> requests = captor.getAllValues();
272+
// This timeout path checks expiration only when new records arrive, so slower CI can shift the
273+
// split points. What should remain stable is that the timeout forces multiple batches while
274+
// preserving every message exactly once.
275+
assertThat(requests.size()).isBetween(2, 3);
276+
assertThat(flattenEntries(requests))
277+
.extracting(SendMessageBatchRequestEntry::messageBody)
278+
.containsExactly("0", "1", "2", "3", "4");
279+
assertThat(requests)
280+
.allSatisfy(
281+
req -> {
282+
assertThat(req.queueUrl()).isEqualTo("queue");
283+
assertThat(req.entries().size()).isBetween(1, 3);
284+
});
270285
}
271286

272287
@Test
@@ -337,11 +352,17 @@ public void testWriteBatchesToDynamicWithTimeout() {
337352

338353
p.run().waitUntilFinish();
339354

340-
SendMessageBatchRequestEntry[] entries = entries(range(0, 5));
341-
// due to added delay, dynamic batches are timed out on arrival of every 2nd msg (per batch)
342-
verify(sqs).sendMessageBatch(request("even", entries[0], entries[2]));
343-
verify(sqs).sendMessageBatch(request("uneven", entries[1], entries[3]));
344-
verify(sqs).sendMessageBatch(request("even", entries[4]));
355+
ArgumentCaptor<SendMessageBatchRequest> captor =
356+
ArgumentCaptor.forClass(SendMessageBatchRequest.class);
357+
verify(sqs, atLeastOnce()).sendMessageBatch(captor.capture());
358+
359+
List<SendMessageBatchRequest> requests = captor.getAllValues();
360+
// Non-strict timeout checks may submit an expired batch either when its next record arrives or
361+
// during a synchronous expiration scan triggered by another queue, so CI timing can vary.
362+
assertThat(requests.size()).isBetween(3, 5);
363+
assertThat(messagesForQueue(requests, "even")).containsExactly("0", "2", "4");
364+
assertThat(messagesForQueue(requests, "uneven")).containsExactly("1", "3");
365+
assertThat(requests).allSatisfy(req -> assertThat(req.entries().size()).isBetween(1, 2));
345366
}
346367

347368
@Test
@@ -406,6 +427,19 @@ private SendMessageBatchRequest anyRequest() {
406427
return any();
407428
}
408429

430+
private List<SendMessageBatchRequestEntry> flattenEntries(
431+
List<SendMessageBatchRequest> requests) {
432+
return requests.stream().flatMap(req -> req.entries().stream()).collect(toList());
433+
}
434+
435+
private List<String> messagesForQueue(List<SendMessageBatchRequest> requests, String queue) {
436+
return requests.stream()
437+
.filter(req -> queue.equals(req.queueUrl()))
438+
.flatMap(req -> req.entries().stream())
439+
.map(SendMessageBatchRequestEntry::messageBody)
440+
.collect(toList());
441+
}
442+
409443
private SendMessageBatchRequest request(String queue, SendMessageBatchRequestEntry... entries) {
410444
return SendMessageBatchRequest.builder()
411445
.queueUrl(queue)

0 commit comments

Comments
 (0)