Apache camel Hazelcast queue polling for concurrency

Viewed 94

I have a requirement for polling a hazelcast (client mode) queue with retry (10 attempts) option on exception. I was expecting that camel polling and processing would be multi threaded. but It wasn't. While retrying on exception, any new message to the queue will be piled up and will be picked up for processing only after 1st one gets completed. Is there any option for parallel processing (concurrent consume). I have added concurrentConsumer and poolSize as a query parameter. But it didn't really play well.

What I have tried is:

fromF(hazelcast-queue://FOO?concurrentConsumers=5&hazelcastInstance=#hazelcastInstance&poolSize=10&queueConsumerMode=Poll).to("direct:testPoll");

from("direct:testPoll")
     .log(LoggingLevel.DEBUG,":::>:Camel[${routeId}] consumes")
     .onException(Exception.class)
     .maximumRedeliveries(maxAttempt)
     .delayPattern(delayPattern)
     .maximumRedeliveryDelay(maxDelay)
     .handled(true)
     .logExhausted(false)
     .end()
.bean("processTestPoll").log(INFO,"${body}").end();

Error:

There are 1 parameters that couldn't be set on the endpoint. Check the uri if the parameters are spelt correctly and that they are properties of the endpoint. Unknown parameters=[{concurrentConsumers=10}]

Your help will be really appreciated. Thanks in advance.

1 Answers

What you try to achieve can be done thanks to a SEDA in 2 different ways:

Generic Way

You can send your messages to a SEDA endpoint and consume them concurrently as next:

fromF("hazelcast-%sFOO?hazelcastInstance=#hazelcastInstance&queueConsumerMode=Poll", 
      HazelcastConstants.QUEUE_PREFIX)
   .to("seda:process");
from("seda:process?concurrentConsumers=5")
   .log("Processing: ${threadName} ${body}");

In the previous example, the Hazelcast Queue FOO is polled by one thread that puts the messages into the SEDA process and the SEDA process is consumed concurrently by 5 threads.

More details about concurrent consumers with the SEDA component

Specific Way

As you proposed in your deleted answer, you can also implement it directly using the specific SEDA endpoint for Hazelcast as next:

fromF("hazelcast-%sFOO?hazelcastInstance=#hazelcastInstance&concurrentConsumers=5", 
      HazelcastConstants.SEDA_PREFIX)
   .log("Processing: ${threadName} ${body}");

In the previous example, the Hazelcast Queue FOO is consumed concurrently by 5 threads.

More details about the Hazelcat SEDA endpoint.

Related