davsclaus commented on code in PR #27243:
URL: https://github.com/apache/camel/pull/27243#discussion_r4163317789
##########
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:
Nit (optional, non-blocking): small pre-existing thing, while this loop is
being restructured anyway. If the thread is interrupted while `isRunAllowed()`
is still true, the interrupt flag is restored and the loop continues. The next
`queue.poll(...)` then throws `InterruptedException` straight away, so the loop
spins hot until the consumer stops. Exiting the loop on interrupt avoids that,
since `doStop` is the only thing that should interrupt this thread (via
`shutdownNow`):
```suggestion
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
```
##########
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:
Nit (optional, non-blocking): if the route stops while this delay is
running, `shutdownNow` interrupts the thread, and `ForegroundTask` then logs a
WARN ("Interrupted HazelcastQueuePollErrorDelay while waiting for the
repeatable task to finish") during a normal shutdown. Harmless, but users may
find it confusing. Also, with `pollingTimeout=0` the error delay is 0, so a
disconnected client would call the exception handler in a tight loop. Would it
be worth using a minimum delay (for example `Math.max(pollingTimeout, 1000)`,
similar to `onErrorDelay` in hazelcast-seda), or checking `isRunAllowed()`
before warning?
--
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]