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]