I have a Spring application that launches 2 consumers which has their own ExecutorServices. In each of the consumer, there I spawn N number of infinite poller loops.
Since theyre SQS pollers, I need a way to gracefully shut them down on a SIGTERM/SIGINT by our process manager. However, it seems like when there are multiple consumers, it seems to shut down one at a time. I want to parallelize this because I dont want the other poller to keep polling while this one is flushing.
class SQSManagerA(
threadPoolSize: Int,
handler: EC2SQSEventHandler,
sqsConsumerEnable: Boolean,
backOffRetry: RetryStrategy<Void>
) {
private val executor: ExecutorService
private var pollers: ArrayList<PollingLoop>
private val waitForTermination = CountDownLatch(1)
private val log = LoggerFactory.getLogger(SQSManagerA::class.java)
init {
executor = Executors.newFixedThreadPool(threadPoolSize)
pollers = ArrayList()
if (sqsConsumerEnable) {
for (i: Int in IntStream.range(0, threadPoolSize)) {
val poller = PollingLoop(handler, backOffRetry)
pollers.add(poller)
executor.execute(poller)
}
}
}
fun destroy() {
log.info("Received a shutdown signal. Signal stop on all threads and wait for 60 seconds")
for (poller in pollers) {
poller.stop()
}
waitForTermination.await(delayElbTargetPoller, TimeUnit.SECONDS)
log.info("Final shutdown on all threads.")
executor.shutdownNow()
}
}
The poller.stop() is just a boolean flag change to let the SQS worker break its infinite loop.
class PollingLoopA(
private val queue: NotificationQueueHandler<>,
private val handler: EventHandler<>,
) : Runnable {
@Volatile
private var terminate = false
override fun run() {
log.info(queue.sourceQueueName + " polling loop is running")
var failedAttempts = 0
while (!terminate) {
//work
}
log.info("${currentThread().name} has stopped polling.")
}
fun stop() {
log.info("${currentThread().name} has called to stop polling.")
this.terminate = true
}
}
The result
[2022-07-23T21:52:57,445] [INFO] (QueueA Polling Thread 28) {} PollingLoopA: No messageA notifications in the queue.
[2022-07-23T21:52:57,475] [INFO] (QueueB Polling Thread 25) {} PollingLoopB: No messageB notifications in the queue.
[2022-07-23T21:53:17,621] [INFO] (QueueB Polling Thread 25) {} PollingLoopB: No messageB notifications in the queue.
[2022-07-23T21:53:17,629] [INFO] (QueueA Polling Thread 28) {} PollingLoopA: No messageA notifications in the queue.
[2022-07-23T21:53:18,413] [INFO] (SIGINT handler) {} spring.Launcher: Initiating shutdown
[2022-07-23T21:53:18,414] [INFO] (SIGINT handler) {} SQSManagerA: Received a shutdown signal. Signal stop on all pollers and wait for 300 seconds
[2022-07-23T21:53:18,415] [INFO] (SIGINT handler) {} PollingLoopA: SIGINT handler has called to stop polling.
[2022-07-23T21:53:37,814] [INFO] (QueueB Polling Thread 25) {} PollingLoopB: No messageB notifications in the queue.
[2022-07-23T21:53:37,828] [INFO] (QueueA Polling Thread 28) {} PollingLoopA: No messageA notifications in the queue.
[2022-07-23T21:53:37,909] [INFO] (QueueA Polling Thread 28) {} PollingLoopA: QueueA Polling Thread 28 has stopped polling.
[2022-07-23T21:53:57,970] [INFO] (QueueB Polling Thread 25) {} PollingLoopB: No messageB notifications in the queue.
[2022-07-23T21:54:18,125] [INFO] (QueueB Polling Thread 25) {} PollingLoopB: No messageB notifications in the queue.
[2022-07-23T21:54:38,273] [INFO] (QueueB Polling Thread 25) {} PollingLoopB: No messageB notifications in the queue.
What happened to my SQSManagerB where it also has a destroyMethod that is the same as SQSManagerA?
Sorry made it generic for NDA purposes
Edit: Op, I found out it actually does terminate both. But it is only one at a time, so if it catches the consumer with the bigger timeout, it waits until then to trigger the other consumer... Is there a way to parallelize the termination?