AWS SQS Listening vs Polling

Viewed 573

I currently have implemented in a Spring Boot project running on Fargate an SQS listener.

It's possible that under the hood, the SqsAsyncClient which appears to be a listener, is actually polling though.

Separately, as a PoC, on I implemented a Lambda function trigger on a different queue. This would be invoked when there are items in the queue and would post to my service. This seems unnecessarily complex to me but removes a single point of failure if I were to only have one instance of the service.

I guess my major point of confusion is whether I am needlessly worrying about polling vs listening on a SQS queue and whether it matters.

Code for example purposes:

@Component
@Slf4j
@RequiredArgsConstructor
public class SqsListener {

private final SqsAsyncClient sqsAsyncClient;
private final Environment environment;
private final SmsMessagingServiceImpl smsMessagingService;

@PostConstruct
public void continuousListener() {
    String queueUrl = environment.getProperty("aws.sqs.sms.queueUrl");
    Mono<ReceiveMessageResponse> responseMono = receiveMessage(queueUrl);
    Flux<Message> messages = getItems(responseMono);
    messages.subscribe(message -> disposeOfFlux(message, queueUrl));
}

protected Flux<Message> getItems(Mono<ReceiveMessageResponse> responseMono) {
   return responseMono.repeat().retry()
            .map(ReceiveMessageResponse::messages)
            .map(Flux::fromIterable)
            .flatMap(messageFlux -> messageFlux);

}

protected void disposeOfFlux(Message message, String queueUrl) {
    log.info("Inbound SMS Received from SQS with MessageId: {}", message.messageId());
    if (someConditionIsMet()) 
        deleteMessage(queueUrl, message);
}

protected Mono<ReceiveMessageResponse> receiveMessage(String queueUrl) {
    return Mono.fromFuture(() -> sqsAsyncClient.receiveMessage(
                    ReceiveMessageRequest.builder()
                            .maxNumberOfMessages(5)
                            .messageAttributeNames("All")
                            .queueUrl(queueUrl)
                            .waitTimeSeconds(10)
                            .visibilityTimeout(30)
                            .build()));
}

protected void deleteMessage(String queueUrl, Message message) {
    sqsAsyncClient.deleteMessage(DeleteMessageRequest.builder()
                    .queueUrl(queueUrl)
                    .receiptHandle(message.receiptHandle())
                    .build())
            .thenAccept(deleteMessageResponse -> log.info("deleted message with handle {}", message.receiptHandle()));
}
}
0 Answers
Related