apache camel-kafka bridgeErrorHandler not working

Viewed 696

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

0 Answers
Related