This is an automated email from the ASF dual-hosted git repository. gnodet pushed a commit to branch fix/CAMEL-24305 in repository https://gitbox.apache.org/repos/asf/camel.git
commit bdcda0caf346666ad083deb006fdb1c568037c6d Author: Guillaume Nodet <[email protected]> AuthorDate: Thu Jul 30 08:39:29 2026 +0200 CAMEL-24305: Fix autowiring of KafkaClientFactory via addComponent() Move default KafkaClientFactory creation from doInit() to doStart() in KafkaComponent. When a component is registered via addComponent() (the path used by Spring Boot), doInit() runs before the autowiring lifecycle strategy, causing the default factory to block injection of a custom one. Also restore the DefaultKafkaClientFactory fallback in KafkaEndpoint's doBuild() method for resilience when the endpoint builds before the component starts. Co-Authored-By: Claude Opus 4.6 <[email protected]> --- .../camel/component/kafka/KafkaComponent.java | 12 +++--- .../camel/component/kafka/KafkaEndpoint.java | 3 ++ .../camel/component/kafka/KafkaAutowireTest.java | 49 +++++++++++++++++++++- 3 files changed, 55 insertions(+), 9 deletions(-) diff --git a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaComponent.java b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaComponent.java index 3883d192bb20..551105bbff68 100644 --- a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaComponent.java +++ b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaComponent.java @@ -298,21 +298,19 @@ public class KafkaComponent extends HealthCheckComponent implements SSLContextPa } @Override - protected void doInit() throws Exception { - super.doInit(); + protected void doStart() throws Exception { + super.doStart(); // if a factory was not autowired then create a default factory + // NOTE: must be done in doStart() rather than doInit(), because when a component is + // registered via addComponent() (the path used by Spring Boot), doInit() runs before + // the autowiring lifecycle strategy has a chance to inject a custom factory. if (kafkaClientFactory == null) { kafkaClientFactory = new DefaultKafkaClientFactory(); } if (configuration.isAllowManualCommit() && kafkaManualCommitFactory == null) { LOG.warn("The component was setup for allowing manual commits, but a manual commit factory was not set"); } - } - - @Override - protected void doStart() throws Exception { - super.doStart(); Map<String, Object> map = new HashMap<>(); // resolve parameter values from the values (#bean / #class etc) diff --git a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaEndpoint.java b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaEndpoint.java index 41549edbba45..30746c1de173 100644 --- a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaEndpoint.java +++ b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaEndpoint.java @@ -132,6 +132,9 @@ public class KafkaEndpoint extends DefaultEndpoint implements MultipleConsumersS if (kafkaClientFactory == null) { kafkaClientFactory = getComponent().getKafkaClientFactory(); } + if (kafkaClientFactory == null) { + kafkaClientFactory = new DefaultKafkaClientFactory(); + } if (kafkaManualCommitFactory == null) { kafkaManualCommitFactory = getComponent().getKafkaManualCommitFactory(); } diff --git a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaAutowireTest.java b/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaAutowireTest.java index d96842acab98..1493d3f96c89 100644 --- a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaAutowireTest.java +++ b/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaAutowireTest.java @@ -18,14 +18,16 @@ package org.apache.camel.component.kafka; import org.apache.camel.BindToRegistry; import org.apache.camel.CamelContext; +import org.apache.camel.impl.DefaultCamelContext; import org.apache.camel.test.infra.core.CamelContextExtension; import org.apache.camel.test.infra.core.DefaultCamelContextExtension; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.RegisterExtension; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertSame; -public class KafkaAutowireTest { +class KafkaAutowireTest { @RegisterExtension protected static CamelContextExtension contextExtension = new DefaultCamelContextExtension(); @@ -36,7 +38,7 @@ public class KafkaAutowireTest { private final KafkaClientFactory clientFactory = new TestKafkaClientFactory(); @Test - public void testKafkaComponentAutowiring() { + void testKafkaComponentAutowiring() { KafkaComponent component = context.getComponent("kafka", KafkaComponent.class); assertSame(clientFactory, component.getKafkaClientFactory()); @@ -44,6 +46,49 @@ public class KafkaAutowireTest { assertSame(clientFactory, endpoint.getKafkaClientFactory()); } + /** + * Verifies that autowiring works when the component is registered via + * {@link CamelContext#addComponent(String, org.apache.camel.Component)}, which is the path used by Spring Boot. In + * this path, the component is added before the context is started, so {@code doInit()} runs before the autowiring + * lifecycle strategy. The default factory must not be created until {@code doStart()} to give autowiring a chance + * to inject a custom factory. + * + * This is a regression test for <a href="https://issues.apache.org/jira/browse/CAMEL-24305">CAMEL-24305</a>. + */ + @Test + void testKafkaComponentAutowiringViaAddComponent() throws Exception { + // Simulate the Spring Boot path: addComponent() is called before the context is started + try (DefaultCamelContext ctx = new DefaultCamelContext()) { + KafkaClientFactory customFactory = new TestKafkaClientFactory(); + ctx.getRegistry().bind("kafkaClientFactory", customFactory); + + KafkaComponent component = new KafkaComponent(); + ctx.addComponent("kafka", component); + + ctx.start(); + + assertSame(customFactory, component.getKafkaClientFactory(), + "Custom KafkaClientFactory should be autowired when using addComponent()"); + } + } + + /** + * Verifies that when no custom KafkaClientFactory is registered, the component creates a + * {@link DefaultKafkaClientFactory} as the default. + */ + @Test + void testKafkaComponentDefaultFactoryWhenNoneRegistered() throws Exception { + try (DefaultCamelContext ctx = new DefaultCamelContext()) { + KafkaComponent component = new KafkaComponent(); + ctx.addComponent("kafka", component); + + ctx.start(); + + assertInstanceOf(DefaultKafkaClientFactory.class, component.getKafkaClientFactory(), + "Default KafkaClientFactory should be created when none is registered"); + } + } + static final class TestKafkaClientFactory extends DefaultKafkaClientFactory { }
