2929import software .amazon .awssdk .services .sqs .model .Message ;
3030import software .amazon .awssdk .services .sqs .model .ReceiveMessageResponse ;
3131
32- import java .util .Collections ;
32+ import java .util .ArrayList ;
3333import java .util .List ;
3434import java .util .concurrent .TimeUnit ;
3535
@@ -44,6 +44,8 @@ public class SqsMessageReceiver {
4444 private static final int SLEEP_DELAY = 20 ;
4545 /** Retry Limit. */
4646 private static final int RETRY_LIMIT = 10 ;
47+ /** Matching message retry limit. */
48+ private static final int MATCHING_MESSAGE_RETRY_LIMIT = 500 ;
4749 /** {@link SqsService}. */
4850 private final SqsService sqs ;
4951 /** Sqs Queue Url. */
@@ -74,8 +76,11 @@ public SqsMessageReceiver(final SqsService sqsService, final String sqsQueueUrl)
7476 * Clears all messages in queue.
7577 */
7678 public void clear () {
77- ReceiveMessageResponse response = sqs .receiveMessages (queueUrl );
78- response .messages ().forEach (m -> sqs .deleteMessage (queueUrl , m .receiptHandle ()));
79+ ReceiveMessageResponse response ;
80+ do {
81+ response = sqs .receiveMessages (queueUrl , MESSAGE_COUNT );
82+ response .messages ().forEach (m -> sqs .deleteMessage (queueUrl , m .receiptHandle ()));
83+ } while (!response .messages ().isEmpty ());
7984 }
8085
8186 /**
@@ -108,23 +113,40 @@ public ReceiveMessageResponse get() throws InterruptedException {
108113 * @throws InterruptedException InterruptedException
109114 */
110115 public List <Message > get (final List <String > searchStrings ) throws InterruptedException {
116+ return get (searchStrings , 1 );
117+ }
118+
119+ /**
120+ * Get the expected number of matching messages.
121+ *
122+ * @param searchStrings {@link List} {@link String}
123+ * @param expectedCount expected number of matching messages
124+ * @return {@link List} {@link Message}
125+ * @throws InterruptedException InterruptedException
126+ */
127+ public List <Message > get (final List <String > searchStrings , final int expectedCount )
128+ throws InterruptedException {
129+ if (expectedCount < 1 ) {
130+ throw new IllegalArgumentException ("'expectedCount' must be greater than zero" );
131+ }
132+
111133 int retry = 0 ;
112- List <Message > messages = Collections . emptyList ();
134+ List <Message > messages = new ArrayList <> ();
113135
114- while (messages .isEmpty () ) {
136+ while (messages .size () < expectedCount ) {
115137 ReceiveMessageResponse response = sqs .receiveMessages (queueUrl , MESSAGE_COUNT );
116- messages = response .messages ();
117- messages .forEach (m -> sqs .deleteMessage (queueUrl , m .receiptHandle ()));
118- messages = messages .stream ().filter (m -> searchStrings .stream ().allMatch (m .body ()::contains ))
119- .toList ();
138+ response .messages ().forEach (m -> sqs .deleteMessage (queueUrl , m .receiptHandle ()));
139+ messages .addAll (response .messages ().stream ()
140+ .filter (m -> searchStrings .stream ().allMatch (m .body ()::contains )).toList ());
120141
121- if (messages .isEmpty () ) {
142+ if (messages .size () < expectedCount ) {
122143 retry ++;
123144 TimeUnit .MILLISECONDS .sleep (SLEEP_DELAY );
124145 }
125146
126- if (retry > RETRY_LIMIT ) {
127- throw new RuntimeException ("Timeout waiting for SQS message" );
147+ if (retry > MATCHING_MESSAGE_RETRY_LIMIT ) {
148+ throw new RuntimeException ("Timeout waiting for " + expectedCount
149+ + " matching SQS messages; received " + messages .size ());
128150 }
129151 }
130152
0 commit comments