allthingssecurity commented on code in PR #27243:
URL: https://github.com/apache/camel/pull/27243#discussion_r4163564380


##########
components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/queue/HazelcastQueueConsumer.java:
##########
@@ -67,41 +81,54 @@ protected void doStop() throws Exception {
 
     class QueueConsumerTask implements Runnable {
 
-        CamelItemListener camelItemListener;
+        private final IQueue<Object> queue;
 
-        public QueueConsumerTask(CamelItemListener camelItemListener) {
-            this.camelItemListener = camelItemListener;
+        QueueConsumerTask(IQueue<Object> queue) {
+            this.queue = queue;
         }
 
         @Override
         public void run() {
-            IQueue<Object> queue = hazelcastInstance.getQueue(cacheName);
-            if (config.getQueueConsumerMode() == 
HazelcastQueueConsumerMode.LISTEN) {
-                queue.addItemListener(camelItemListener, true);
-            }
-
-            if (config.getQueueConsumerMode() == 
HazelcastQueueConsumerMode.POLL) {
-                while (isRunAllowed()) {
+            while (isRunAllowed()) {
+                final Object body;
+                try {
+                    body = queue.poll(config.getPollingTimeout(), 
TimeUnit.MILLISECONDS);
+                } catch (InterruptedException e) {
+                    Thread.currentThread().interrupt();
+                    continue;

Review Comment:
   Applied in 238b1ecc8a1f: on `InterruptedException` the loop restores the 
flag and returns. The error-delay wait also returns from `run()` if it is 
interrupted.
   
   _Claude Code on behalf of allthingssecurity_



##########
components/camel-hazelcast/src/main/java/org/apache/camel/component/hazelcast/queue/HazelcastQueueConsumer.java:
##########
@@ -67,41 +81,54 @@ protected void doStop() throws Exception {
 
     class QueueConsumerTask implements Runnable {
 
-        CamelItemListener camelItemListener;
+        private final IQueue<Object> queue;
 
-        public QueueConsumerTask(CamelItemListener camelItemListener) {
-            this.camelItemListener = camelItemListener;
+        QueueConsumerTask(IQueue<Object> queue) {
+            this.queue = queue;
         }
 
         @Override
         public void run() {
-            IQueue<Object> queue = hazelcastInstance.getQueue(cacheName);
-            if (config.getQueueConsumerMode() == 
HazelcastQueueConsumerMode.LISTEN) {
-                queue.addItemListener(camelItemListener, true);
-            }
-
-            if (config.getQueueConsumerMode() == 
HazelcastQueueConsumerMode.POLL) {
-                while (isRunAllowed()) {
+            while (isRunAllowed()) {
+                final Object body;
+                try {
+                    body = queue.poll(config.getPollingTimeout(), 
TimeUnit.MILLISECONDS);
+                } catch (InterruptedException e) {
+                    Thread.currentThread().interrupt();
+                    continue;
+                } catch (Exception e) {
+                    // keep polling after an error (such as the client being 
disconnected from the cluster)
+                    if (isRunAllowed()) {
+                        getExceptionHandler().handleException("Error polling 
from the queue " + cacheName, e);
+                        waitBeforeNextPoll();
+                    }
+                    continue;
+                }
+                // CAMEL-16035 - If the polling timeout is exceeded with 
nothing to poll from the queue, the queue.poll() method return NULL
+                if (body != null) {
+                    Exchange exchange = createExchange(false);
+                    exchange.getIn().setBody(body);
                     try {
-                        final Object body = 
queue.poll(config.getPollingTimeout(), TimeUnit.MILLISECONDS);
-                        // CAMEL-16035 - If the polling timeout is exceeded 
with nothing to poll from the queue, the queue.poll() method return NULL
-                        if (body != null) {
-                            Exchange exchange = createExchange(false);
-                            exchange.getIn().setBody(body);
-                            try {
-                                processor.process(exchange);
-                            } catch (Exception e) {
-                                getExceptionHandler().handleException("Error 
during processing", exchange, e);
-                            } finally {
-                                releaseExchange(exchange, false);
-                            }
-                        }
-                    } catch (InterruptedException e) {
-                        Thread.currentThread().interrupt();
+                        processor.process(exchange);
+                    } catch (Exception e) {
+                        getExceptionHandler().handleException("Error during 
processing", exchange, e);
+                    } finally {
+                        releaseExchange(exchange, false);
                     }
                 }
             }
         }
+
+        private void waitBeforeNextPoll() {

Review Comment:
   Done in 238b1ecc8a1f. The delay is now `Math.max(pollingTimeout, 1000)`. The 
WARN is logged inside `ForegroundTask` itself, so checking `isRunAllowed()` on 
our side couldn't avoid it. Instead the wait is a `CountDownLatch.await(delay)` 
on a latch that `doStop` counts down before `shutdownNow`. A stop ends the 
delay at once and nothing is logged. All camel-hazelcast tests pass (237, the 
same unrelated `HazelcastAggregationRepositoryRoutesTest` flake as on main).
   
   _Claude Code on behalf of allthingssecurity_



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to