I am trying to handle the exceptions that occurred while the consumer is trying to pick up incoming messages, such as deserialization error.As I read the document there is a property bridgeErrorHandlersetting this to true enables the exception/error handler configured in the route.
Below is the camel route configuration
try (CamelContext camelContext = new DefaultCamelContext()) {
camelContext.addRoutes(new RouteBuilder() {
@Override
public void configure() throws Exception {
onCompletion().process(exchange -> {
KafkaManualCommit manual = exchange.getIn().getHeader(KafkaConstants.MANUAL_COMMIT,
KafkaManualCommit.class);
if (null != manual) {
manual.commitSync();
}
});
onException(Exception.class)
.maximumRedeliveries(2)
.redeliveryDelay(2)
.handled(true)
.process(e -> {
Exception ex = e.getProperty(Exchange.EXCEPTION_CAUGHT, Exception.class);
log.error("error occured on route, will be handled",ex);
//eror handling logic
});
from("kafka:test123456?brokers=localhost:9092&allowManualCommit=true"
+ "&autoCommitEnable=false&valueDeserializer=org.apache.kafka.common.serialization.LongDeserializer"
+ "&bridgeErrorHandler=true&breakOnFirstError=false")
.routeId("route-123")
.process(e->{
e.getIn().setBody("enriched-message");
}).id("process-123").description("enrich-processor")
.to("stream:out").id("sink-123");
}
});
camelContext.start();
Howe ever, exception handler is not triggered for any KafkaException and it is ignored and logged by the camel default error handler. Is there anything missing in the configuration? I am using camel-kafka 3.5.0 version