This is an automated email from the ASF dual-hosted git repository.
Croway pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new c01088d815ed CAMEL-25029: camel-kafka - move the shared client code to
camel-kafka-common (#27488)
c01088d815ed is described below
commit c01088d815edd871560c61dccb07b0ebadabf55e
Author: Federico Mariani <[email protected]>
AuthorDate: Wed Oct 7 16:19:43 2026 +0200
CAMEL-25029: camel-kafka - move the shared client code to
camel-kafka-common (#27488)
* CAMEL-25029: camel-kafka - move the shared client code to
camel-kafka-common
The kafka-share consumer must be built on the Kafka client code without
depending on camel-kafka, so the shared classes move to a new flat
camel-kafka-common module, with their packages and names unchanged:
KafkaClientConfiguration, AbstractKafkaComponent, KafkaConstants,
KafkaHeaderFilterStrategy, PollExceptionStrategy, PollOnError,
KafkaConsumerFatalException, TaskHealthState, KafkaRecordProcessor,
JMSDeserializer, the serde and the security packages.
camel-kafka depends on camel-kafka-common, so applications and runtimes
that depend on camel-kafka need no change. KafkaClientFactory stays in
camel-kafka.
The options of the moved classes carry an explicit description, as the
catalog generator cannot read javadoc from a dependency jar; kafka.json
is unchanged.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
* CAMEL-25029: camel-kafka - fix the pollExceptionStrategy and
additionalProperties descriptions
Review comments: "pooling" is "polling", and the additionalProperties
example lost its square brackets. The catalog generator drops [ and ]
from every description, so the example is described in words; the
javadoc keeps the literal example. The upgrade guide lists every class
that moved to camel-kafka-common.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---------
Co-authored-by: Claude Opus 5.5 <[email protected]>
---
.github/actions/labeler/label-config.yml | 1 +
bom/camel-bom/pom.xml | 5 +
.../org/apache/camel/catalog/components/kafka.json | 6 +-
components/camel-kafka-common/pom.xml | 99 ++++++++++
.../component/kafka/AbstractKafkaComponent.java | 23 ++-
.../component/kafka/KafkaClientConfiguration.java | 201 ++++++++++++++++-----
.../camel/component/kafka/KafkaConstants.java | 0
.../kafka/KafkaConsumerFatalException.java | 0
.../component/kafka/KafkaHeaderFilterStrategy.java | 0
.../component/kafka/PollExceptionStrategy.java | 0
.../apache/camel/component/kafka/PollOnError.java | 0
.../camel/component/kafka/TaskHealthState.java | 0
.../consumer/support/KafkaRecordProcessor.java | 0
.../consumer/support/interop/JMSDeserializer.java | 0
.../component/kafka/security/KafkaAuthType.java | 0
.../kafka/security/KafkaSecurityConfigurer.java | 0
.../serde/DefaultKafkaHeaderDeserializer.java | 0
.../kafka/serde/DefaultKafkaHeaderSerializer.java | 0
.../kafka/serde/KafkaHeaderDeserializer.java | 0
.../kafka/serde/KafkaHeaderSerializer.java | 0
.../component/kafka/serde/KafkaSerdeHelper.java | 0
.../serde/ToStringKafkaHeaderDeserializer.java | 0
.../serde/DefaultKafkaHeaderDeserializerTest.java | 0
.../serde/DefaultKafkaHeaderSerializerTest.java | 0
.../serde/ToStringKafkaHeaderDeserializerTest.java | 0
components/camel-kafka/pom.xml | 22 +--
.../org/apache/camel/component/kafka/kafka.json | 6 +-
components/pom.xml | 1 +
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 22 +++
.../dsl/KafkaComponentBuilderFactory.java | 6 +-
.../endpoint/dsl/KafkaEndpointBuilderFactory.java | 24 +--
parent/pom.xml | 5 +
32 files changed, 332 insertions(+), 89 deletions(-)
diff --git a/.github/actions/labeler/label-config.yml
b/.github/actions/labeler/label-config.yml
index baeddb005580..2480344b7cf0 100644
--- a/.github/actions/labeler/label-config.yml
+++ b/.github/actions/labeler/label-config.yml
@@ -82,6 +82,7 @@ components-kafka:
- changed-files:
- any-glob-to-any-file:
- components/camel-kafka/**/*
+ - components/camel-kafka-common/**/*
components-jms:
- changed-files:
diff --git a/bom/camel-bom/pom.xml b/bom/camel-bom/pom.xml
index b8775ddf05af..18bd5e5dfc26 100644
--- a/bom/camel-bom/pom.xml
+++ b/bom/camel-bom/pom.xml
@@ -1512,6 +1512,11 @@
<artifactId>camel-kafka</artifactId>
<version>4.23.0-SNAPSHOT</version>
</dependency>
+ <dependency>
+ <groupId>org.apache.camel</groupId>
+ <artifactId>camel-kafka-common</artifactId>
+ <version>4.23.0-SNAPSHOT</version>
+ </dependency>
<dependency>
<groupId>org.apache.camel</groupId>
<artifactId>camel-kamelet</artifactId>
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/kafka.json
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/kafka.json
index 3e43c941c1e4..0abdfd66bbdd 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/kafka.json
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/kafka.json
@@ -24,7 +24,7 @@
"remote": true
},
"componentProperties": {
- "additionalProperties": { "index": 0, "kind": "property", "displayName":
"Additional Properties", "group": "common", "label": "common", "required":
false, "type": "object", "javaType": "java.util.Map<java.lang.String,
java.lang.Object>", "prefix": "additionalProperties.", "multiValue": true,
"deprecated": false, "autowired": false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Sets [...]
+ "additionalProperties": { "index": 0, "kind": "property", "displayName":
"Additional Properties", "group": "common", "label": "common", "required":
false, "type": "object", "javaType": "java.util.Map<java.lang.String,
java.lang.Object>", "prefix": "additionalProperties.", "multiValue": true,
"deprecated": false, "autowired": false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Sets [...]
"brokers": { "index": 1, "kind": "property", "displayName": "Brokers",
"group": "common", "label": "common", "required": false, "type": "string",
"javaType": "java.lang.String", "deprecated": false, "autowired": false,
"secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "URL of the Kafka brokers to use. The format is
host1:port1,host2:port2, and the list can be a subset of brokers or a VIP [...]
"clientId": { "index": 2, "kind": "property", "displayName": "Client Id",
"group": "common", "label": "common", "required": false, "type": "string",
"javaType": "java.lang.String", "deprecated": false, "autowired": false,
"secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "The client id is a user-specified string sent
in each request to help trace calls. It should logically identify the a [...]
"configuration": { "index": 3, "kind": "property", "displayName":
"Configuration", "group": "common", "label": "", "required": false, "type":
"object", "javaType": "org.apache.camel.component.kafka.KafkaConfiguration",
"deprecated": false, "autowired": false, "secret": false, "description":
"Allows to pre-configure the Kafka component with common options that the
endpoints will reuse." },
@@ -80,7 +80,7 @@
"createConsumerBackoffMaxAttempts": { "index": 53, "kind": "property",
"displayName": "Create Consumer Backoff Max Attempts", "group": "consumer
(advanced)", "label": "consumer,advanced", "required": false, "type":
"integer", "javaType": "int", "deprecated": false, "autowired": false,
"secret": false, "description": "Maximum attempts to create the kafka consumer
(kafka-client), before eventually giving up and failing. Error during creating
the consumer may be fatal due to invalid con [...]
"isolationLevel": { "index": 54, "kind": "property", "displayName":
"Isolation Level", "group": "consumer (advanced)", "label":
"consumer,advanced", "required": false, "type": "enum", "javaType":
"java.lang.String", "enum": [ "read_uncommitted", "read_committed" ],
"deprecated": false, "autowired": false, "secret": false, "defaultValue":
"read_uncommitted", "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description [...]
"kafkaManualCommitFactory": { "index": 55, "kind": "property",
"displayName": "Kafka Manual Commit Factory", "group": "consumer (advanced)",
"label": "consumer,advanced", "required": false, "type": "object", "javaType":
"org.apache.camel.component.kafka.consumer.KafkaManualCommitFactory",
"deprecated": false, "autowired": true, "secret": false, "description":
"Factory to use for creating KafkaManualCommit instances. This allows to plugin
a custom factory to create custom KafkaManualC [...]
- "pollExceptionStrategy": { "index": 56, "kind": "property", "displayName":
"Poll Exception Strategy", "group": "consumer (advanced)", "label":
"consumer,advanced", "required": false, "type": "object", "javaType":
"org.apache.camel.component.kafka.PollExceptionStrategy", "deprecated": false,
"autowired": true, "secret": false, "description": "To use a custom strategy
with the consumer to control how to handle exceptions thrown from the Kafka
broker while pooling messages." },
+ "pollExceptionStrategy": { "index": 56, "kind": "property", "displayName":
"Poll Exception Strategy", "group": "consumer (advanced)", "label":
"consumer,advanced", "required": false, "type": "object", "javaType":
"org.apache.camel.component.kafka.PollExceptionStrategy", "deprecated": false,
"autowired": true, "secret": false, "description": "To use a custom strategy
with the consumer to control how to handle exceptions thrown from the Kafka
broker while polling messages." },
"subscribeConsumerBackoffInterval": { "index": 57, "kind": "property",
"displayName": "Subscribe Consumer Backoff Interval", "group": "consumer
(advanced)", "label": "consumer,advanced", "required": false, "type":
"integer", "javaType": "long", "deprecated": true, "autowired": false,
"secret": false, "defaultValue": 5000, "description": "The delay in millis
seconds to wait before trying again to subscribe to the kafka broker." },
"subscribeConsumerBackoffMaxAttempts": { "index": 58, "kind": "property",
"displayName": "Subscribe Consumer Backoff Max Attempts", "group": "consumer
(advanced)", "label": "consumer,advanced", "required": false, "type":
"integer", "javaType": "int", "deprecated": true, "autowired": false, "secret":
false, "description": "Maximum number the kafka consumer will attempt to
subscribe to the kafka broker, before eventually giving up and failing. Error
during subscribing the consumer to t [...]
"subscribeConsumerTopicMustExists": { "index": 59, "kind": "property",
"displayName": "Subscribe Consumer Topic Must Exists", "group": "consumer
(advanced)", "label": "consumer,advanced", "required": false, "type":
"boolean", "javaType": "boolean", "deprecated": false, "autowired": false,
"secret": false, "defaultValue": false, "description": "Whether when a Camel
Kafka consumer is subscribing to a Kafka broker then check whether a topic
already exist on the broker, and fail if it do [...]
@@ -171,7 +171,7 @@
},
"properties": {
"topic": { "index": 0, "kind": "path", "displayName": "Topic", "group":
"common", "label": "common", "required": true, "type": "string", "javaType":
"java.lang.String", "deprecated": false, "deprecationNote": "", "autowired":
false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Name of the topic to use. On the consumer you
can use comma to separate multiple topics. A producer can on [...]
- "additionalProperties": { "index": 1, "kind": "parameter", "displayName":
"Additional Properties", "group": "common", "label": "common", "required":
false, "type": "object", "javaType": "java.util.Map<java.lang.String,
java.lang.Object>", "prefix": "additionalProperties.", "multiValue": true,
"deprecated": false, "autowired": false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Sets [...]
+ "additionalProperties": { "index": 1, "kind": "parameter", "displayName":
"Additional Properties", "group": "common", "label": "common", "required":
false, "type": "object", "javaType": "java.util.Map<java.lang.String,
java.lang.Object>", "prefix": "additionalProperties.", "multiValue": true,
"deprecated": false, "autowired": false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Sets [...]
"brokers": { "index": 2, "kind": "parameter", "displayName": "Brokers",
"group": "common", "label": "common", "required": false, "type": "string",
"javaType": "java.lang.String", "deprecated": false, "autowired": false,
"secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "URL of the Kafka brokers to use. The format is
host1:port1,host2:port2, and the list can be a subset of brokers or a VI [...]
"clientId": { "index": 3, "kind": "parameter", "displayName": "Client Id",
"group": "common", "label": "common", "required": false, "type": "string",
"javaType": "java.lang.String", "deprecated": false, "autowired": false,
"secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "The client id is a user-specified string sent
in each request to help trace calls. It should logically identify the [...]
"connectionMaxIdleMs": { "index": 4, "kind": "parameter", "displayName":
"Connection Max Idle Ms", "group": "common", "label": "common", "required":
false, "type": "integer", "javaType": "java.lang.Integer", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": 540000,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "Close idle connections
after the number of milliseconds specified [...]
diff --git a/components/camel-kafka-common/pom.xml
b/components/camel-kafka-common/pom.xml
new file mode 100644
index 000000000000..9d2324b4dcbe
--- /dev/null
+++ b/components/camel-kafka-common/pom.xml
@@ -0,0 +1,99 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+
+ Licensed to the Apache Software Foundation (ASF) under one or more
+ contributor license agreements. See the NOTICE file distributed with
+ this work for additional information regarding copyright ownership.
+ The ASF licenses this file to You under the Apache License, Version 2.0
+ (the "License"); you may not use this file except in compliance with
+ the License. You may obtain a copy of the License at
+
+ http://www.apache.org/licenses/LICENSE-2.0
+
+ Unless required by applicable law or agreed to in writing, software
+ distributed under the License is distributed on an "AS IS" BASIS,
+ WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ See the License for the specific language governing permissions and
+ limitations under the License.
+
+-->
+<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/maven-v4_0_0.xsd">
+ <modelVersion>4.0.0</modelVersion>
+
+ <parent>
+ <groupId>org.apache.camel</groupId>
+ <artifactId>components</artifactId>
+ <version>4.23.0-SNAPSHOT</version>
+ </parent>
+
+ <artifactId>camel-kafka-common</artifactId>
+ <packaging>jar</packaging>
+ <name>Camel :: Kafka :: Common</name>
+ <description>Camel Kafka common shared code</description>
+
+ <dependencies>
+
+ <!-- camel -->
+ <dependency>
+ <groupId>org.apache.camel</groupId>
+ <artifactId>camel-support</artifactId>
+ </dependency>
+
+ <!-- kafka java client -->
+ <dependency>
+ <groupId>org.apache.kafka</groupId>
+ <artifactId>kafka-clients</artifactId>
+ <version>${kafka-version}</version>
+ <exclusions>
+ <exclusion>
+ <groupId>org.lz4</groupId>
+ <artifactId>lz4-java</artifactId>
+ </exclusion>
+ </exclusions>
+ </dependency>
+ <dependency>
+ <groupId>at.yawk.lz4</groupId>
+ <artifactId>lz4-java</artifactId>
+ <version>${lz4-java-version}</version>
+ </dependency>
+
+ <!-- test -->
+ <dependency>
+ <groupId>org.apache.camel</groupId>
+ <artifactId>camel-jackson</artifactId>
+ <scope>test</scope>
+ </dependency>
+ <dependency>
+ <groupId>org.hamcrest</groupId>
+ <artifactId>hamcrest</artifactId>
+ <version>${hamcrest-version}</version>
+ <scope>test</scope>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.camel</groupId>
+ <artifactId>camel-test-junit6</artifactId>
+ <scope>test</scope>
+ </dependency>
+ </dependencies>
+
+ <build>
+ <plugins>
+ <!-- This module is a shared library, not a Camel component -->
+ <plugin>
+ <groupId>org.apache.camel</groupId>
+ <artifactId>camel-package-maven-plugin</artifactId>
+ <executions>
+ <execution>
+ <id>generate</id>
+ <phase>none</phase>
+ </execution>
+ <execution>
+ <id>generate-postcompile</id>
+ <phase>none</phase>
+ </execution>
+ </executions>
+ </plugin>
+ </plugins>
+ </build>
+
+</project>
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/AbstractKafkaComponent.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/AbstractKafkaComponent.java
similarity index 76%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/AbstractKafkaComponent.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/AbstractKafkaComponent.java
index 59a44d0410e4..f390c183cb5e 100644
---
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/AbstractKafkaComponent.java
+++
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/AbstractKafkaComponent.java
@@ -39,13 +39,26 @@ public abstract class AbstractKafkaComponent extends
HealthCheckComponent
private final List<Runnable> pendingConsumers = new
CopyOnWriteArrayList<>();
- @Metadata(label = "security", defaultValue = "false")
+ @Metadata(label = "security", defaultValue = "false", description =
"Enable usage of global SSL context parameters.")
private boolean useGlobalSslContextParameters;
- @Metadata(autowired = true, label = "consumer,advanced")
+ @Metadata(autowired = true, label = "consumer,advanced",
+ description = "To use a custom strategy with the consumer to
control how to handle exceptions thrown from the "
+ + "Kafka broker while polling messages.")
private PollExceptionStrategy pollExceptionStrategy;
- @Metadata(label = "consumer,advanced")
+ @Metadata(label = "consumer,advanced",
+ description = "Maximum attempts to create the kafka consumer
(kafka-client), before eventually giving up and "
+ + "failing. Error during creating the consumer may
be fatal due to invalid configuration and as "
+ + "such recovery is not possible. However, one
part of the validation is DNS resolution of the "
+ + "bootstrap broker hostnames. This may be a
temporary networking problem, and could potentially "
+ + "be recoverable. While other errors are fatal,
such as some invalid kafka configurations. "
+ + "Unfortunately, kafka-client does not separate
this kind of errors. Camel will by default retry "
+ + "forever, and therefore never give up. If you
want to give up after many attempts then set this "
+ + "option and Camel will then when giving up
terminate the consumer. To try again, you can "
+ + "manually restart the consumer by stopping, and
starting the route.")
private int createConsumerBackoffMaxAttempts;
- @Metadata(label = "consumer,advanced", defaultValue = "5000")
+ @Metadata(label = "consumer,advanced", defaultValue = "5000",
+ description = "The delay in millis seconds to wait before trying
again to create the kafka consumer "
+ + "(kafka-client).")
private long createConsumerBackoffInterval = 5000;
protected AbstractKafkaComponent() {
@@ -78,7 +91,7 @@ public abstract class AbstractKafkaComponent extends
HealthCheckComponent
/**
* To use a custom strategy with the consumer to control how to handle
exceptions thrown from the Kafka broker while
- * pooling messages.
+ * polling messages.
*/
public void setPollExceptionStrategy(PollExceptionStrategy
pollExceptionStrategy) {
this.pollExceptionStrategy = pollExceptionStrategy;
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaClientConfiguration.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/KafkaClientConfiguration.java
similarity index 76%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaClientConfiguration.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/KafkaClientConfiguration.java
index 1e7fc1dcef1c..e310bcab71fe 100644
---
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaClientConfiguration.java
+++
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/KafkaClientConfiguration.java
@@ -62,129 +62,225 @@ import
org.apache.kafka.common.security.auth.SecurityProtocol;
@UriParams
public abstract class KafkaClientConfiguration implements Cloneable,
HeaderFilterStrategyAware {
- @UriParam(label = "common")
+ @UriParam(label = "common",
+ description = "URL of the Kafka brokers to use. The format is
host1:port1,host2:port2, and the list can be a "
+ + "subset of brokers or a VIP pointing to a subset
of brokers. This option is known as "
+ + "bootstrap.servers in the Kafka documentation.")
private String brokers;
- @UriParam(label = "common")
+ @UriParam(label = "common",
+ description = "The client id is a user-specified string sent in
each request to help trace calls. It should "
+ + "logically identify the application making the
request.")
private String clientId;
@UriParam(label = "common",
description = "To use a custom HeaderFilterStrategy to filter
header to and from Camel message.")
private HeaderFilterStrategy headerFilterStrategy = new
KafkaHeaderFilterStrategy();
- @UriParam(label = "common", defaultValue = "100")
+ @UriParam(label = "common", defaultValue = "100",
+ description = "The amount of time to wait before attempting to
retry a failed request to a given topic "
+ + "partition. This avoids repeatedly sending
requests in a tight loop under some failure "
+ + "scenarios. This value is the initial backoff
value and will increase exponentially for each "
+ + "failed request, up to the retry.backoff.max.ms
value.")
private Integer retryBackoffMs = 100;
- @UriParam(label = "common", defaultValue = "1000")
+ @UriParam(label = "common", defaultValue = "1000",
+ description = "The maximum amount of time in milliseconds to
wait when retrying a request to the broker that "
+ + "has repeatedly failed. If provided, the backoff
per client will increase exponentially for each"
+ + " failed request, up to this maximum. To prevent
all clients from being synchronized upon retry,"
+ + " a randomized jitter with a factor of 0.2 will
be applied to the backoff, resulting in the "
+ + "backoff falling within a range between 20%
below and 20% above the computed value. If "
+ + "retry.backoff.ms is set to be higher than
retry.backoff.max.ms, then retry.backoff.max.ms will "
+ + "be used as a constant backoff from the
beginning without any exponential increase")
private Integer retryBackoffMaxMs = 1000;
- @UriParam(label = "consumer", defaultValue = "true")
+ @UriParam(label = "consumer", defaultValue = "true",
+ description = "Whether to eager validate that broker host:port
is valid and can be DNS resolved to known host "
+ + "during starting this consumer. If the
validation fails, then an exception is thrown, which "
+ + "makes Camel fail fast. Disabling this will
postpone the validation after the consumer is "
+ + "started, and Camel will keep re-connecting in
case of validation or DNS resolution error.")
private boolean preValidateHostAndPort = true;
@UriParam(label = "consumer", description = "To use a custom
KafkaHeaderDeserializer to deserialize kafka headers values")
private KafkaHeaderDeserializer headerDeserializer = new
DefaultKafkaHeaderDeserializer();
// key.deserializer
- @UriParam(label = "consumer", defaultValue =
KafkaConstants.KAFKA_DEFAULT_DESERIALIZER)
+ @UriParam(label = "consumer", defaultValue =
KafkaConstants.KAFKA_DEFAULT_DESERIALIZER,
+ description = "Deserializer class for the key that implements
the Deserializer interface.")
private String keyDeserializer = KafkaConstants.KAFKA_DEFAULT_DESERIALIZER;
// value.deserializer
- @UriParam(label = "consumer", defaultValue =
KafkaConstants.KAFKA_DEFAULT_DESERIALIZER)
+ @UriParam(label = "consumer", defaultValue =
KafkaConstants.KAFKA_DEFAULT_DESERIALIZER,
+ description = "Deserializer class for value that implements the
Deserializer interface.")
private String valueDeserializer =
KafkaConstants.KAFKA_DEFAULT_DESERIALIZER;
// connections.max.idle.ms
- @UriParam(label = "common", defaultValue = "540000")
+ @UriParam(label = "common", defaultValue = "540000",
+ description = "Close idle connections after the number of
milliseconds specified by this config.")
private Integer connectionMaxIdleMs = 540000;
// receive.buffer.bytes
- @UriParam(label = "common", defaultValue = "65536")
+ @UriParam(label = "common", defaultValue = "65536",
+ description = "The size of the TCP receive buffer (SO_RCVBUF) to
use when reading data.")
private Integer receiveBufferBytes = 65536;
// send.buffer.bytes
- @UriParam(label = "common", defaultValue = "131072")
+ @UriParam(label = "common", defaultValue = "131072", description = "Socket
write buffer size")
private Integer sendBufferBytes = 131072;
// metadata.max.age.ms
- @UriParam(label = "common", defaultValue = "300000")
+ @UriParam(label = "common", defaultValue = "300000",
+ description = "The period of time in milliseconds after which we
force a refresh of metadata even if we "
+ + "haven't seen any partition leadership changes
to proactively discover any new brokers or "
+ + "partitions.")
private Integer metadataMaxAgeMs = 300000;
// metric.reporters
- @UriParam(label = "common")
+ @UriParam(label = "common",
+ description = "A list of classes to use as metrics reporters.
Implementing the MetricReporter interface allows"
+ + " plugging in classes that will be notified of
new metric creation. The JmxReporter is always "
+ + "included to register JMX statistics.")
private String metricReporters;
// metrics.num.samples
- @UriParam(label = "common", defaultValue = "2")
+ @UriParam(label = "common", defaultValue = "2", description = "The number
of samples maintained to compute metrics.")
private Integer noOfMetricsSample = 2;
// metrics.sample.window.ms
- @UriParam(label = "common", defaultValue = "30000")
+ @UriParam(label = "common", defaultValue = "30000", description = "The
window of time a metrics sample is computed over.")
private Integer metricsSampleWindowMs = 30000;
// reconnect.backoff.ms
- @UriParam(label = "common", defaultValue = "50")
+ @UriParam(label = "common", defaultValue = "50",
+ description = "The amount of time to wait before attempting to
reconnect to a given host. This avoids "
+ + "repeatedly connecting to a host in a tight
loop. This backoff applies to all requests sent by "
+ + "the consumer to the broker.")
private Integer reconnectBackoffMs = 50;
// reconnect.backoff.max.ms
- @UriParam(label = "common", defaultValue = "1000")
+ @UriParam(label = "common", defaultValue = "1000",
+ description = "The maximum amount of time in milliseconds to
wait when reconnecting to a broker that has "
+ + "repeatedly failed to connect. If provided, the
backoff per host will increase exponentially for"
+ + " each consecutive connection failure, up to
this maximum. After calculating the backoff "
+ + "increase, 20% random jitter is added to avoid
connection storms.")
private Integer reconnectBackoffMaxMs = 1000;
// SSL
- @UriParam(label = "common,security")
+ @UriParam(label = "common,security",
+ description = "SSL configuration using a Camel
SSLContextParameters object. If configured, it's applied before"
+ + " the other SSL endpoint parameters. NOTE: Kafka
only supports loading keystore from file "
+ + "locations, so prefix the location with file: in
the KeyStoreParameters.resource option.")
private SSLContextParameters sslContextParameters;
// SSL
// ssl.key.password
- @UriParam(label = "common,security", security = "secret")
+ @UriParam(label = "common,security", security = "secret",
+ description = "The password of the private key in the key store
file or the PEM key specified in "
+ + "sslKeystoreKey. This is required for clients
only if two-way authentication is configured.")
private String sslKeyPassword;
// ssl.keystore.location
- @UriParam(label = "common,security")
+ @UriParam(label = "common,security",
+ description = "The location of the key store file. This is
optional for the client and can be used for two-way"
+ + " authentication for the client.")
private String sslKeystoreLocation;
// ssl.keystore.password
- @UriParam(label = "common,security", security = "secret")
+ @UriParam(label = "common,security", security = "secret",
+ description = "The store password for the key store file. This
is optional for the client and only needed if "
+ + "sslKeystoreLocation is configured. Key store
password is not supported for PEM format.")
private String sslKeystorePassword;
// ssl.truststore.location
- @UriParam(label = "common,security")
+ @UriParam(label = "common,security", description = "The location of the
trust store file.")
private String sslTruststoreLocation;
// ssl.truststore.password
- @UriParam(label = "common,security", security = "secret")
+ @UriParam(label = "common,security", security = "secret",
+ description = "The password for the trust store file. If a
password is not set, trust store file configured "
+ + "will still be used, but integrity checking is
disabled. Trust store password is not supported "
+ + "for PEM format.")
private String sslTruststorePassword;
// SSL
// ssl.enabled.protocols
- @UriParam(label = "common,security")
+ @UriParam(label = "common,security",
+ description = "The list of protocols enabled for SSL
connections. The default is TLSv1.2,TLSv1.3 when running "
+ + "with Java 11 or newer, TLSv1.2 otherwise. With
the default value for Java 11, clients and "
+ + "servers will prefer TLSv1.3 if both support it
and fallback to TLSv1.2 otherwise (assuming both"
+ + " support at least TLSv1.2). This default should
be fine for most cases. Also see the config "
+ + "documentation for SslProtocol.")
private String sslEnabledProtocols =
SslConfigs.DEFAULT_SSL_ENABLED_PROTOCOLS;
// ssl.keystore.type
- @UriParam(label = "common,security", defaultValue =
SslConfigs.DEFAULT_SSL_KEYSTORE_TYPE)
+ @UriParam(label = "common,security", defaultValue =
SslConfigs.DEFAULT_SSL_KEYSTORE_TYPE,
+ description = "The file format of the key store file. This is
optional for the client. The default value is "
+ + "JKS")
private String sslKeystoreType = SslConfigs.DEFAULT_SSL_KEYSTORE_TYPE;
// ssl.protocol
- @UriParam(label = "common,security")
+ @UriParam(label = "common,security",
+ description = "The SSL protocol used to generate the SSLContext.
The default is TLSv1.3 when running with Java"
+ + " 11 or newer, TLSv1.2 otherwise. This value
should be fine for most use cases. Allowed values "
+ + "in recent JVMs are TLSv1.2 and TLSv1.3. TLS,
TLSv1.1, SSL, SSLv2 and SSLv3 may be supported in "
+ + "older JVMs, but their usage is discouraged due
to known security vulnerabilities. With the "
+ + "default value for this config and
sslEnabledProtocols, clients will downgrade to TLSv1.2 if the"
+ + " server does not support TLSv1.3. If this
config is set to TLSv1.2, clients will not use "
+ + "TLSv1.3 even if it is one of the values in
sslEnabledProtocols and the server only supports "
+ + "TLSv1.3.")
private String sslProtocol = SslConfigs.DEFAULT_SSL_PROTOCOL;
// ssl.provider
- @UriParam(label = "common,security")
+ @UriParam(label = "common,security",
+ description = "The name of the security provider used for SSL
connections. Default value is the default "
+ + "security provider of the JVM.")
private String sslProvider;
// ssl.truststore.type
- @UriParam(label = "common,security", defaultValue =
SslConfigs.DEFAULT_SSL_TRUSTSTORE_TYPE)
+ @UriParam(label = "common,security", defaultValue =
SslConfigs.DEFAULT_SSL_TRUSTSTORE_TYPE,
+ description = "The file format of the trust store file. The
default value is JKS.")
private String sslTruststoreType = SslConfigs.DEFAULT_SSL_TRUSTSTORE_TYPE;
// SSL
// ssl.cipher.suites
- @UriParam(label = "common,security")
+ @UriParam(label = "common,security",
+ description = "A list of cipher suites. This is a named
combination of authentication, encryption, MAC and key"
+ + " exchange algorithm used to negotiate the
security settings for a network connection using TLS "
+ + "or SSL network protocol. By default, all the
available cipher suites are supported.")
private String sslCipherSuites;
// ssl.endpoint.identification.algorithm
- @UriParam(label = "common,security", defaultValue = "https", security =
"insecure:ssl", insecureValue = "none")
+ @UriParam(label = "common,security", defaultValue = "https", security =
"insecure:ssl", insecureValue = "none",
+ description = "The endpoint identification algorithm to validate
server hostname using server certificate. Use"
+ + " none or false to disable server hostname
verification.")
private String sslEndpointAlgorithm =
SslConfigs.DEFAULT_SSL_ENDPOINT_IDENTIFICATION_ALGORITHM;
// ssl.keymanager.algorithm
- @UriParam(label = "common,security", defaultValue = "SunX509")
+ @UriParam(label = "common,security", defaultValue = "SunX509",
+ description = "The algorithm used by key manager factory for SSL
connections. Default value is the key manager"
+ + " factory algorithm configured for the Java
Virtual Machine.")
private String sslKeymanagerAlgorithm = "SunX509";
// ssl.trustmanager.algorithm
- @UriParam(label = "common,security", defaultValue = "PKIX")
+ @UriParam(label = "common,security", defaultValue = "PKIX",
+ description = "The algorithm used by trust manager factory for
SSL connections. Default value is the trust "
+ + "manager factory algorithm configured for the
Java Virtual Machine.")
private String sslTrustmanagerAlgorithm = "PKIX";
// SASL & sucurity Protocol
// sasl.kerberos.service.name
- @UriParam(label = "common,security")
+ @UriParam(label = "common,security",
+ description = "The Kerberos principal name that Kafka runs as.
This can be defined either in Kafka's JAAS "
+ + "config or in Kafka's config.")
private String saslKerberosServiceName;
// security.protocol
- @UriParam(label = "common,security", defaultValue =
CommonClientConfigs.DEFAULT_SECURITY_PROTOCOL)
+ @UriParam(label = "common,security", defaultValue =
CommonClientConfigs.DEFAULT_SECURITY_PROTOCOL,
+ description = "Protocol used to communicate with brokers.
SASL_PLAINTEXT, PLAINTEXT, SASL_SSL and SSL are "
+ + "supported")
private String securityProtocol =
CommonClientConfigs.DEFAULT_SECURITY_PROTOCOL;
// SASL
// sasl.mechanism
- @UriParam(label = "common,security", defaultValue =
SaslConfigs.DEFAULT_SASL_MECHANISM)
+ @UriParam(label = "common,security", defaultValue =
SaslConfigs.DEFAULT_SASL_MECHANISM,
+ description = "The Simple Authentication and Security Layer
(SASL) Mechanism used. For the valid values see "
+ +
"http://www.iana.org/assignments/sasl-mechanisms/sasl-mechanisms.xhtml")
private String saslMechanism = SaslConfigs.DEFAULT_SASL_MECHANISM;
// sasl.kerberos.kinit.cmd
- @UriParam(label = "common,security", defaultValue =
SaslConfigs.DEFAULT_KERBEROS_KINIT_CMD)
+ @UriParam(label = "common,security", defaultValue =
SaslConfigs.DEFAULT_KERBEROS_KINIT_CMD,
+ description = "Kerberos kinit command path. Default is
/usr/bin/kinit")
private String kerberosInitCmd = SaslConfigs.DEFAULT_KERBEROS_KINIT_CMD;
// sasl.kerberos.min.time.before.relogin
- @UriParam(label = "common,security", defaultValue = "60000")
+ @UriParam(label = "common,security", defaultValue = "60000",
+ description = "Login thread sleep time between refresh
attempts.")
private Integer kerberosBeforeReloginMinTime = 60000;
// sasl.kerberos.ticket.renew.jitter
- @UriParam(label = "common,security", defaultValue = "0.05")
+ @UriParam(label = "common,security", defaultValue = "0.05",
+ description = "Percentage of random jitter added to the renewal
time.")
private Double kerberosRenewJitter =
SaslConfigs.DEFAULT_KERBEROS_TICKET_RENEW_JITTER;
// sasl.kerberos.ticket.renew.window.factor
- @UriParam(label = "common,security", defaultValue = "0.8")
+ @UriParam(label = "common,security", defaultValue = "0.8",
+ description = "Login thread will sleep until the specified
window factor of time from last refresh to ticket's"
+ + " expiry has been reached, at which time it will
try to renew the ticket.")
private Double kerberosRenewWindowFactor =
SaslConfigs.DEFAULT_KERBEROS_TICKET_RENEW_WINDOW_FACTOR;
- @UriParam(label = "common,security", defaultValue = "DEFAULT")
+ @UriParam(label = "common,security", defaultValue = "DEFAULT",
+ description = "A list of rules for mapping from principal names
to short names (typically operating system "
+ + "usernames). The rules are evaluated in order,
and the first rule that matches a principal name "
+ + "is used to map it to a short name. Any later
rules in the list are ignored. By default, "
+ + "principal names of the form
{username}/{hostname}{REALM} are mapped to {username}. For more "
+ + "details on the format, please see the Security
Authorization and ACLs documentation (at the "
+ + "Apache Kafka project website). Multiple values
can be separated by comma")
// sasl.kerberos.principal.to.local.rules
private String kerberosPrincipalToLocalRules;
- @UriParam(label = "common,security", security = "secret")
+ @UriParam(label = "common,security", security = "secret",
+ description = "Expose the kafka sasl.jaas.config parameter
Example: "
+ +
"org.apache.kafka.common.security.plain.PlainLoginModule required
username=USERNAME "
+ + "password=PASSWORD;")
// sasl.jaas.config
private String saslJaasConfig;
// Simplified authentication configuration
@@ -215,16 +311,31 @@ public abstract class KafkaClientConfiguration implements
Cloneable, HeaderFilte
description = "OAuth scope. Used when saslAuthType is set to
OAUTH.")
private String oauthScope;
// Schema registry only options
- @UriParam(label = "schema")
+ @UriParam(label = "schema",
+ description = "URL of the schema registry servers to use. The
format is host1:port1,host2:port2. This is known"
+ + " as schema.registry.url in multiple Schema
registries documentation. This option is only "
+ + "available externally (not standard Apache
Kafka)")
private String schemaRegistryURL;
- @UriParam(label = "schema,consumer")
+ @UriParam(label = "schema,consumer",
+ description = "This enables the use of a specific Avro reader
for use with the in multiple Schema registries "
+ + "documentation with Avro Deserializers
implementation. This option is only available externally "
+ + "(not standard Apache Kafka)")
private boolean specificAvroReader;
// Additional properties
- @UriParam(label = "common", prefix = "additionalProperties.", multiValue =
true)
+ @UriParam(label = "common", prefix = "additionalProperties.", multiValue =
true,
+ description = "Sets additional properties for either kafka
consumer or kafka producer in case they can't be "
+ + "set directly on the camel configurations (e.g.:
new Kafka properties that are not reflected yet"
+ + " in Camel configurations), the properties have
to be prefixed with additionalProperties., e.g.:"
+ + "
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://lo"
+ + "calhost:8811/avro. If the properties are set in
the application.properties file, they must be "
+ + "prefixed with
camel.component.kafka.additional-properties followed by the property name "
+ + "enclosed in square brackets, for example the
delivery.timeout.ms property in square brackets.")
private Map<String, Object> additionalProperties = new HashMap<>();
- @UriParam(label = "common", defaultValue = "30000")
+ @UriParam(label = "common", defaultValue = "30000",
+ description = "Timeout in milliseconds to wait gracefully for
the consumer or producer to shut down and "
+ + "terminate its worker threads.")
private int shutdownTimeout = 30000;
- @UriParam(label = "common,security")
+ @UriParam(label = "common,security", description = "Location of the
kerberos config file.")
private String kerberosConfigLocation;
/**
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaConstants.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/KafkaConstants.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaConstants.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/KafkaConstants.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaConsumerFatalException.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/KafkaConsumerFatalException.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaConsumerFatalException.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/KafkaConsumerFatalException.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaHeaderFilterStrategy.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/KafkaHeaderFilterStrategy.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaHeaderFilterStrategy.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/KafkaHeaderFilterStrategy.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/PollExceptionStrategy.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/PollExceptionStrategy.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/PollExceptionStrategy.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/PollExceptionStrategy.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/PollOnError.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/PollOnError.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/PollOnError.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/PollOnError.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/TaskHealthState.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/TaskHealthState.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/TaskHealthState.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/TaskHealthState.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/support/KafkaRecordProcessor.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/consumer/support/KafkaRecordProcessor.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/support/KafkaRecordProcessor.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/consumer/support/KafkaRecordProcessor.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/support/interop/JMSDeserializer.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/consumer/support/interop/JMSDeserializer.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/support/interop/JMSDeserializer.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/consumer/support/interop/JMSDeserializer.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/security/KafkaAuthType.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/security/KafkaAuthType.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/security/KafkaAuthType.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/security/KafkaAuthType.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/security/KafkaSecurityConfigurer.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/security/KafkaSecurityConfigurer.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/security/KafkaSecurityConfigurer.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/security/KafkaSecurityConfigurer.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderDeserializer.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderDeserializer.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderDeserializer.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderDeserializer.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderSerializer.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderSerializer.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderSerializer.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderSerializer.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/serde/KafkaHeaderDeserializer.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/serde/KafkaHeaderDeserializer.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/serde/KafkaHeaderDeserializer.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/serde/KafkaHeaderDeserializer.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/serde/KafkaHeaderSerializer.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/serde/KafkaHeaderSerializer.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/serde/KafkaHeaderSerializer.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/serde/KafkaHeaderSerializer.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/serde/KafkaSerdeHelper.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/serde/KafkaSerdeHelper.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/serde/KafkaSerdeHelper.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/serde/KafkaSerdeHelper.java
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/serde/ToStringKafkaHeaderDeserializer.java
b/components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/serde/ToStringKafkaHeaderDeserializer.java
similarity index 100%
rename from
components/camel-kafka/src/main/java/org/apache/camel/component/kafka/serde/ToStringKafkaHeaderDeserializer.java
rename to
components/camel-kafka-common/src/main/java/org/apache/camel/component/kafka/serde/ToStringKafkaHeaderDeserializer.java
diff --git
a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderDeserializerTest.java
b/components/camel-kafka-common/src/test/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderDeserializerTest.java
similarity index 100%
rename from
components/camel-kafka/src/test/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderDeserializerTest.java
rename to
components/camel-kafka-common/src/test/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderDeserializerTest.java
diff --git
a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderSerializerTest.java
b/components/camel-kafka-common/src/test/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderSerializerTest.java
similarity index 100%
rename from
components/camel-kafka/src/test/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderSerializerTest.java
rename to
components/camel-kafka-common/src/test/java/org/apache/camel/component/kafka/serde/DefaultKafkaHeaderSerializerTest.java
diff --git
a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/serde/ToStringKafkaHeaderDeserializerTest.java
b/components/camel-kafka-common/src/test/java/org/apache/camel/component/kafka/serde/ToStringKafkaHeaderDeserializerTest.java
similarity index 100%
rename from
components/camel-kafka/src/test/java/org/apache/camel/component/kafka/serde/ToStringKafkaHeaderDeserializerTest.java
rename to
components/camel-kafka-common/src/test/java/org/apache/camel/component/kafka/serde/ToStringKafkaHeaderDeserializerTest.java
diff --git a/components/camel-kafka/pom.xml b/components/camel-kafka/pom.xml
index b2ad7bb4f01c..c08fbfc35593 100644
--- a/components/camel-kafka/pom.xml
+++ b/components/camel-kafka/pom.xml
@@ -34,6 +34,10 @@
<dependencies>
<!-- camel -->
+ <dependency>
+ <groupId>org.apache.camel</groupId>
+ <artifactId>camel-kafka-common</artifactId>
+ </dependency>
<dependency>
<groupId>org.apache.camel</groupId>
<artifactId>camel-support</artifactId>
@@ -43,24 +47,6 @@
<artifactId>camel-health</artifactId>
</dependency>
- <!-- kafka java client -->
- <dependency>
- <groupId>org.apache.kafka</groupId>
- <artifactId>kafka-clients</artifactId>
- <version>${kafka-version}</version>
- <exclusions>
- <exclusion>
- <groupId>org.lz4</groupId>
- <artifactId>lz4-java</artifactId>
- </exclusion>
- </exclusions>
- </dependency>
- <dependency>
- <groupId>at.yawk.lz4</groupId>
- <artifactId>lz4-java</artifactId>
- <version>${lz4-java-version}</version>
- </dependency>
-
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
diff --git
a/components/camel-kafka/src/generated/resources/META-INF/org/apache/camel/component/kafka/kafka.json
b/components/camel-kafka/src/generated/resources/META-INF/org/apache/camel/component/kafka/kafka.json
index 3e43c941c1e4..0abdfd66bbdd 100644
---
a/components/camel-kafka/src/generated/resources/META-INF/org/apache/camel/component/kafka/kafka.json
+++
b/components/camel-kafka/src/generated/resources/META-INF/org/apache/camel/component/kafka/kafka.json
@@ -24,7 +24,7 @@
"remote": true
},
"componentProperties": {
- "additionalProperties": { "index": 0, "kind": "property", "displayName":
"Additional Properties", "group": "common", "label": "common", "required":
false, "type": "object", "javaType": "java.util.Map<java.lang.String,
java.lang.Object>", "prefix": "additionalProperties.", "multiValue": true,
"deprecated": false, "autowired": false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Sets [...]
+ "additionalProperties": { "index": 0, "kind": "property", "displayName":
"Additional Properties", "group": "common", "label": "common", "required":
false, "type": "object", "javaType": "java.util.Map<java.lang.String,
java.lang.Object>", "prefix": "additionalProperties.", "multiValue": true,
"deprecated": false, "autowired": false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Sets [...]
"brokers": { "index": 1, "kind": "property", "displayName": "Brokers",
"group": "common", "label": "common", "required": false, "type": "string",
"javaType": "java.lang.String", "deprecated": false, "autowired": false,
"secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "URL of the Kafka brokers to use. The format is
host1:port1,host2:port2, and the list can be a subset of brokers or a VIP [...]
"clientId": { "index": 2, "kind": "property", "displayName": "Client Id",
"group": "common", "label": "common", "required": false, "type": "string",
"javaType": "java.lang.String", "deprecated": false, "autowired": false,
"secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "The client id is a user-specified string sent
in each request to help trace calls. It should logically identify the a [...]
"configuration": { "index": 3, "kind": "property", "displayName":
"Configuration", "group": "common", "label": "", "required": false, "type":
"object", "javaType": "org.apache.camel.component.kafka.KafkaConfiguration",
"deprecated": false, "autowired": false, "secret": false, "description":
"Allows to pre-configure the Kafka component with common options that the
endpoints will reuse." },
@@ -80,7 +80,7 @@
"createConsumerBackoffMaxAttempts": { "index": 53, "kind": "property",
"displayName": "Create Consumer Backoff Max Attempts", "group": "consumer
(advanced)", "label": "consumer,advanced", "required": false, "type":
"integer", "javaType": "int", "deprecated": false, "autowired": false,
"secret": false, "description": "Maximum attempts to create the kafka consumer
(kafka-client), before eventually giving up and failing. Error during creating
the consumer may be fatal due to invalid con [...]
"isolationLevel": { "index": 54, "kind": "property", "displayName":
"Isolation Level", "group": "consumer (advanced)", "label":
"consumer,advanced", "required": false, "type": "enum", "javaType":
"java.lang.String", "enum": [ "read_uncommitted", "read_committed" ],
"deprecated": false, "autowired": false, "secret": false, "defaultValue":
"read_uncommitted", "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description [...]
"kafkaManualCommitFactory": { "index": 55, "kind": "property",
"displayName": "Kafka Manual Commit Factory", "group": "consumer (advanced)",
"label": "consumer,advanced", "required": false, "type": "object", "javaType":
"org.apache.camel.component.kafka.consumer.KafkaManualCommitFactory",
"deprecated": false, "autowired": true, "secret": false, "description":
"Factory to use for creating KafkaManualCommit instances. This allows to plugin
a custom factory to create custom KafkaManualC [...]
- "pollExceptionStrategy": { "index": 56, "kind": "property", "displayName":
"Poll Exception Strategy", "group": "consumer (advanced)", "label":
"consumer,advanced", "required": false, "type": "object", "javaType":
"org.apache.camel.component.kafka.PollExceptionStrategy", "deprecated": false,
"autowired": true, "secret": false, "description": "To use a custom strategy
with the consumer to control how to handle exceptions thrown from the Kafka
broker while pooling messages." },
+ "pollExceptionStrategy": { "index": 56, "kind": "property", "displayName":
"Poll Exception Strategy", "group": "consumer (advanced)", "label":
"consumer,advanced", "required": false, "type": "object", "javaType":
"org.apache.camel.component.kafka.PollExceptionStrategy", "deprecated": false,
"autowired": true, "secret": false, "description": "To use a custom strategy
with the consumer to control how to handle exceptions thrown from the Kafka
broker while polling messages." },
"subscribeConsumerBackoffInterval": { "index": 57, "kind": "property",
"displayName": "Subscribe Consumer Backoff Interval", "group": "consumer
(advanced)", "label": "consumer,advanced", "required": false, "type":
"integer", "javaType": "long", "deprecated": true, "autowired": false,
"secret": false, "defaultValue": 5000, "description": "The delay in millis
seconds to wait before trying again to subscribe to the kafka broker." },
"subscribeConsumerBackoffMaxAttempts": { "index": 58, "kind": "property",
"displayName": "Subscribe Consumer Backoff Max Attempts", "group": "consumer
(advanced)", "label": "consumer,advanced", "required": false, "type":
"integer", "javaType": "int", "deprecated": true, "autowired": false, "secret":
false, "description": "Maximum number the kafka consumer will attempt to
subscribe to the kafka broker, before eventually giving up and failing. Error
during subscribing the consumer to t [...]
"subscribeConsumerTopicMustExists": { "index": 59, "kind": "property",
"displayName": "Subscribe Consumer Topic Must Exists", "group": "consumer
(advanced)", "label": "consumer,advanced", "required": false, "type":
"boolean", "javaType": "boolean", "deprecated": false, "autowired": false,
"secret": false, "defaultValue": false, "description": "Whether when a Camel
Kafka consumer is subscribing to a Kafka broker then check whether a topic
already exist on the broker, and fail if it do [...]
@@ -171,7 +171,7 @@
},
"properties": {
"topic": { "index": 0, "kind": "path", "displayName": "Topic", "group":
"common", "label": "common", "required": true, "type": "string", "javaType":
"java.lang.String", "deprecated": false, "deprecationNote": "", "autowired":
false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Name of the topic to use. On the consumer you
can use comma to separate multiple topics. A producer can on [...]
- "additionalProperties": { "index": 1, "kind": "parameter", "displayName":
"Additional Properties", "group": "common", "label": "common", "required":
false, "type": "object", "javaType": "java.util.Map<java.lang.String,
java.lang.Object>", "prefix": "additionalProperties.", "multiValue": true,
"deprecated": false, "autowired": false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Sets [...]
+ "additionalProperties": { "index": 1, "kind": "parameter", "displayName":
"Additional Properties", "group": "common", "label": "common", "required":
false, "type": "object", "javaType": "java.util.Map<java.lang.String,
java.lang.Object>", "prefix": "additionalProperties.", "multiValue": true,
"deprecated": false, "autowired": false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Sets [...]
"brokers": { "index": 2, "kind": "parameter", "displayName": "Brokers",
"group": "common", "label": "common", "required": false, "type": "string",
"javaType": "java.lang.String", "deprecated": false, "autowired": false,
"secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "URL of the Kafka brokers to use. The format is
host1:port1,host2:port2, and the list can be a subset of brokers or a VI [...]
"clientId": { "index": 3, "kind": "parameter", "displayName": "Client Id",
"group": "common", "label": "common", "required": false, "type": "string",
"javaType": "java.lang.String", "deprecated": false, "autowired": false,
"secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "The client id is a user-specified string sent
in each request to help trace calls. It should logically identify the [...]
"connectionMaxIdleMs": { "index": 4, "kind": "parameter", "displayName":
"Connection Max Idle Ms", "group": "common", "label": "common", "required":
false, "type": "integer", "javaType": "java.lang.Integer", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": 540000,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "Close idle connections
after the number of milliseconds specified [...]
diff --git a/components/pom.xml b/components/pom.xml
index cc525e0c5f37..30c50bef77da 100644
--- a/components/pom.xml
+++ b/components/pom.xml
@@ -205,6 +205,7 @@
<module>camel-jt400</module>
<module>camel-jta</module>
<module>camel-jte</module>
+ <module>camel-kafka-common</module>
<module>camel-kafka</module>
<module>camel-kamelet</module>
<module>camel-keycloak</module>
diff --git
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index 639c5fdd6ee7..0874e9de60df 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -4185,6 +4185,28 @@ consumed those partitions while the circuit was open;
now they stay with this co
it runs. A suspended consumer also keeps its fetcher thread after a failed
poll or a reconnect, and the health check
reports it as recoverable.
+=== camel-kafka - shared client code moved to camel-kafka-common
+
+The Kafka client code that the `kafka` component shares with other Kafka based
components has been moved
+from `camel-kafka` into a new `camel-kafka-common` module. This includes
`KafkaClientConfiguration` (the base class
+of `KafkaConfiguration`), `AbstractKafkaComponent`, `KafkaConstants`,
`KafkaHeaderFilterStrategy`, the header
+serializers and deserializers (`org.apache.camel.component.kafka.serde`),
`KafkaSecurityConfigurer` and
+`KafkaAuthType`, `PollExceptionStrategy`, `PollOnError`,
`KafkaConsumerFatalException`, `TaskHealthState`,
+`consumer.support.KafkaRecordProcessor` and
`consumer.support.interop.JMSDeserializer`. The packages and class names
+are unchanged.
+
+`camel-kafka` depends on `camel-kafka-common`, so nothing changes for
applications that depend on `camel-kafka`.
+If you build your classpath without transitive dependencies, add the
`camel-kafka-common` dependency:
+
+[source,xml]
+----
+<dependency>
+ <groupId>org.apache.camel</groupId>
+ <artifactId>camel-kafka-common</artifactId>
+ <version>${camel.version}</version>
+</dependency>
+----
+
=== camel-pulsar - the producer no longer replaces the body with the message id
After sending, the producer used to overwrite the message body with the
`MessageId` returned by the
diff --git
a/dsl/camel-componentdsl/src/generated/java/org/apache/camel/builder/component/dsl/KafkaComponentBuilderFactory.java
b/dsl/camel-componentdsl/src/generated/java/org/apache/camel/builder/component/dsl/KafkaComponentBuilderFactory.java
index b8e606a64fe2..1aefe08cf364 100644
---
a/dsl/camel-componentdsl/src/generated/java/org/apache/camel/builder/component/dsl/KafkaComponentBuilderFactory.java
+++
b/dsl/camel-componentdsl/src/generated/java/org/apache/camel/builder/component/dsl/KafkaComponentBuilderFactory.java
@@ -55,8 +55,8 @@ public interface KafkaComponentBuilderFactory {
* producer in case they can't be set directly on the camel
* configurations (e.g.: new Kafka properties that are not reflected
yet
* in Camel configurations), the properties have to be prefixed with
- * additionalProperties.., e.g.:
- *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties and the property
enclosed in square brackets, like this example:
camel.component.kafka.additional-propertiesdelivery.timeout.ms=15000. This is a
multi-value option with prefix: additionalProperties.
+ * additionalProperties., e.g.:
+ *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties followed by the
property name enclosed in square brackets, for example the delivery.timeout.ms
property in square brackets. This is a multi-value option with prefix:
additionalProperties.
*
* The option is a: <code>java.util.Map&lt;java.lang.String,
* java.lang.Object&gt;</code> type.
@@ -1183,7 +1183,7 @@ public interface KafkaComponentBuilderFactory {
/**
* To use a custom strategy with the consumer to control how to handle
- * exceptions thrown from the Kafka broker while pooling messages.
+ * exceptions thrown from the Kafka broker while polling messages.
*
* The option is a:
*
<code>org.apache.camel.component.kafka.PollExceptionStrategy</code>
type.
diff --git
a/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/KafkaEndpointBuilderFactory.java
b/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/KafkaEndpointBuilderFactory.java
index 5a440eb83ece..b44cfdce7b63 100644
---
a/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/KafkaEndpointBuilderFactory.java
+++
b/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/KafkaEndpointBuilderFactory.java
@@ -48,8 +48,8 @@ public interface KafkaEndpointBuilderFactory {
* producer in case they can't be set directly on the camel
* configurations (e.g.: new Kafka properties that are not reflected
yet
* in Camel configurations), the properties have to be prefixed with
- * additionalProperties.., e.g.:
- *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties and the property
enclosed in square brackets, like this example:
camel.component.kafka.additional-propertiesdelivery.timeout.ms=15000. This is a
multi-value option with prefix: additionalProperties.
+ * additionalProperties., e.g.:
+ *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties followed by the
property name enclosed in square brackets, for example the delivery.timeout.ms
property in square brackets. This is a multi-value option with prefix:
additionalProperties.
*
* The option is a: <code>java.util.Map<java.lang.String,
* java.lang.Object></code> type.
@@ -72,8 +72,8 @@ public interface KafkaEndpointBuilderFactory {
* producer in case they can't be set directly on the camel
* configurations (e.g.: new Kafka properties that are not reflected
yet
* in Camel configurations), the properties have to be prefixed with
- * additionalProperties.., e.g.:
- *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties and the property
enclosed in square brackets, like this example:
camel.component.kafka.additional-propertiesdelivery.timeout.ms=15000. This is a
multi-value option with prefix: additionalProperties.
+ * additionalProperties., e.g.:
+ *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties followed by the
property name enclosed in square brackets, for example the delivery.timeout.ms
property in square brackets. This is a multi-value option with prefix:
additionalProperties.
*
* The option is a: <code>java.util.Map<java.lang.String,
* java.lang.Object></code> type.
@@ -2613,8 +2613,8 @@ public interface KafkaEndpointBuilderFactory {
* producer in case they can't be set directly on the camel
* configurations (e.g.: new Kafka properties that are not reflected
yet
* in Camel configurations), the properties have to be prefixed with
- * additionalProperties.., e.g.:
- *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties and the property
enclosed in square brackets, like this example:
camel.component.kafka.additional-propertiesdelivery.timeout.ms=15000. This is a
multi-value option with prefix: additionalProperties.
+ * additionalProperties., e.g.:
+ *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties followed by the
property name enclosed in square brackets, for example the delivery.timeout.ms
property in square brackets. This is a multi-value option with prefix:
additionalProperties.
*
* The option is a: <code>java.util.Map<java.lang.String,
* java.lang.Object></code> type.
@@ -2637,8 +2637,8 @@ public interface KafkaEndpointBuilderFactory {
* producer in case they can't be set directly on the camel
* configurations (e.g.: new Kafka properties that are not reflected
yet
* in Camel configurations), the properties have to be prefixed with
- * additionalProperties.., e.g.:
- *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties and the property
enclosed in square brackets, like this example:
camel.component.kafka.additional-propertiesdelivery.timeout.ms=15000. This is a
multi-value option with prefix: additionalProperties.
+ * additionalProperties., e.g.:
+ *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties followed by the
property name enclosed in square brackets, for example the delivery.timeout.ms
property in square brackets. This is a multi-value option with prefix:
additionalProperties.
*
* The option is a: <code>java.util.Map<java.lang.String,
* java.lang.Object></code> type.
@@ -4910,8 +4910,8 @@ public interface KafkaEndpointBuilderFactory {
* producer in case they can't be set directly on the camel
* configurations (e.g.: new Kafka properties that are not reflected
yet
* in Camel configurations), the properties have to be prefixed with
- * additionalProperties.., e.g.:
- *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties and the property
enclosed in square brackets, like this example:
camel.component.kafka.additional-propertiesdelivery.timeout.ms=15000. This is a
multi-value option with prefix: additionalProperties.
+ * additionalProperties., e.g.:
+ *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties followed by the
property name enclosed in square brackets, for example the delivery.timeout.ms
property in square brackets. This is a multi-value option with prefix:
additionalProperties.
*
* The option is a: <code>java.util.Map<java.lang.String,
* java.lang.Object></code> type.
@@ -4934,8 +4934,8 @@ public interface KafkaEndpointBuilderFactory {
* producer in case they can't be set directly on the camel
* configurations (e.g.: new Kafka properties that are not reflected
yet
* in Camel configurations), the properties have to be prefixed with
- * additionalProperties.., e.g.:
- *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties and the property
enclosed in square brackets, like this example:
camel.component.kafka.additional-propertiesdelivery.timeout.ms=15000. This is a
multi-value option with prefix: additionalProperties.
+ * additionalProperties., e.g.:
+ *
additionalProperties.transactional.id=12345&additionalProperties.schema.registry.url=http://localhost:8811/avro.
If the properties are set in the application.properties file, they must be
prefixed with camel.component.kafka.additional-properties followed by the
property name enclosed in square brackets, for example the delivery.timeout.ms
property in square brackets. This is a multi-value option with prefix:
additionalProperties.
*
* The option is a: <code>java.util.Map<java.lang.String,
* java.lang.Object></code> type.
diff --git a/parent/pom.xml b/parent/pom.xml
index ae9e47c5a13a..584813d9007f 100644
--- a/parent/pom.xml
+++ b/parent/pom.xml
@@ -2055,6 +2055,11 @@
<artifactId>camel-kafka</artifactId>
<version>${project.version}</version>
</dependency>
+ <dependency>
+ <groupId>org.apache.camel</groupId>
+ <artifactId>camel-kafka-common</artifactId>
+ <version>${project.version}</version>
+ </dependency>
<dependency>
<groupId>org.apache.camel</groupId>
<artifactId>camel-kamelet</artifactId>