This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch camel-4.22.x
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/camel-4.22.x by this push:
new 337a87d552d6 CAMEL-25091: camel-kafka - expose by_duration as valid
autoOffsetReset value (#27271)
337a87d552d6 is described below
commit 337a87d552d6750d861cec5ee0045e57686064a0
Author: Salvatore Mongiardo <[email protected]>
AuthorDate: Fri Oct 2 12:16:25 2026 +0200
CAMEL-25091: camel-kafka - expose by_duration as valid autoOffsetReset
value (#27271)
Backport of #26988 to camel-4.22.x. Kafka 4.0 added the
by_duration:<ISO-8601> auto-offset-reset strategy; autoOffsetReset no longer
restricts the value to an enum so it passes catalog and camel validate checks,
and its description is corrected.
---
.../org/apache/camel/catalog/components/kafka.json | 4 +-
.../apache/camel/catalog/docs/kafka-component.adoc | 43 ++++++++++++++++++++++
.../org/apache/camel/component/kafka/kafka.json | 4 +-
.../camel-kafka/src/main/docs/kafka-component.adoc | 43 ++++++++++++++++++++++
.../camel/component/kafka/KafkaConfiguration.java | 9 +++--
.../component/kafka/KafkaConfigurationTest.java | 17 +++++++++
.../dsl/KafkaComponentBuilderFactory.java | 10 +++--
.../endpoint/dsl/KafkaEndpointBuilderFactory.java | 10 +++--
.../tui/PropertyCompletionProviderTest.java | 12 +++---
.../core/commands/tui/YamlCompletionTest.java | 15 ++++++--
10 files changed, 142 insertions(+), 25 deletions(-)
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 6317731965ad..c5eda1d3b834 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
@@ -37,7 +37,7 @@
"allowManualCommit": { "index": 10, "kind": "property", "displayName":
"Allow Manual Commit", "group": "consumer", "label": "consumer", "required":
false, "type": "boolean", "javaType": "boolean", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": false,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "Whether to allow doing
manual commits via KafkaManualCommit. If this option is [...]
"autoCommitEnable": { "index": 11, "kind": "property", "displayName":
"Auto Commit Enable", "group": "consumer", "label": "consumer", "required":
false, "type": "boolean", "javaType": "boolean", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": true,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "If true, periodically
commit to ZooKeeper the offset of messages already fetched [...]
"autoCommitIntervalMs": { "index": 12, "kind": "property", "displayName":
"Auto Commit Interval Ms", "group": "consumer", "label": "consumer",
"required": false, "type": "integer", "javaType": "java.lang.Integer",
"deprecated": false, "autowired": false, "secret": false, "defaultValue": 5000,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "The frequency in ms that
the consumer offsets are committed to [...]
- "autoOffsetReset": { "index": 13, "kind": "property", "displayName": "Auto
Offset Reset", "group": "consumer", "label": "consumer", "required": false,
"type": "enum", "javaType": "java.lang.String", "enum": [ "latest", "earliest",
"none" ], "deprecated": false, "autowired": false, "secret": false,
"defaultValue": "latest", "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "What to do when there is no ini [...]
+ "autoOffsetReset": { "index": 13, "kind": "property", "displayName": "Auto
Offset Reset", "group": "consumer", "label": "consumer", "required": false,
"type": "string", "javaType": "java.lang.String", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": "latest",
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "Where a consumer group
starts reading when it has no committed offset, [...]
"batching": { "index": 14, "kind": "property", "displayName": "Batching",
"group": "consumer", "label": "consumer", "required": false, "type": "boolean",
"javaType": "boolean", "deprecated": false, "autowired": false, "secret":
false, "defaultValue": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Whether to use batching for processing or
streaming. The default is false, which uses streaming. I [...]
"batchingIntervalMs": { "index": 15, "kind": "property", "displayName":
"Batching Interval Ms", "group": "consumer", "label": "consumer", "required":
false, "type": "integer", "javaType": "java.lang.Integer", "deprecated": false,
"autowired": false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "In consumer batching mode, then this option is
specifying a time in millis, to trigger ba [...]
"breakOnFirstError": { "index": 16, "kind": "property", "displayName":
"Break On First Error", "group": "consumer", "label": "consumer", "required":
false, "type": "boolean", "javaType": "boolean", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": false,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "This options controls
what happens when a consumer is processing an exchange [...]
@@ -183,7 +183,7 @@
"allowManualCommit": { "index": 10, "kind": "parameter", "displayName":
"Allow Manual Commit", "group": "consumer", "label": "consumer", "required":
false, "type": "boolean", "javaType": "boolean", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": false,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "Whether to allow doing
manual commits via KafkaManualCommit. If this option i [...]
"autoCommitEnable": { "index": 11, "kind": "parameter", "displayName":
"Auto Commit Enable", "group": "consumer", "label": "consumer", "required":
false, "type": "boolean", "javaType": "boolean", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": true,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "If true, periodically
commit to ZooKeeper the offset of messages already fetched [...]
"autoCommitIntervalMs": { "index": 12, "kind": "parameter", "displayName":
"Auto Commit Interval Ms", "group": "consumer", "label": "consumer",
"required": false, "type": "integer", "javaType": "java.lang.Integer",
"deprecated": false, "autowired": false, "secret": false, "defaultValue": 5000,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "The frequency in ms that
the consumer offsets are committed t [...]
- "autoOffsetReset": { "index": 13, "kind": "parameter", "displayName":
"Auto Offset Reset", "group": "consumer", "label": "consumer", "required":
false, "type": "enum", "javaType": "java.lang.String", "enum": [ "latest",
"earliest", "none" ], "deprecated": false, "autowired": false, "secret": false,
"defaultValue": "latest", "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "What to do when there is no in [...]
+ "autoOffsetReset": { "index": 13, "kind": "parameter", "displayName":
"Auto Offset Reset", "group": "consumer", "label": "consumer", "required":
false, "type": "string", "javaType": "java.lang.String", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": "latest",
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "Where a consumer group
starts reading when it has no committed offset, [...]
"batching": { "index": 14, "kind": "parameter", "displayName": "Batching",
"group": "consumer", "label": "consumer", "required": false, "type": "boolean",
"javaType": "boolean", "deprecated": false, "autowired": false, "secret":
false, "defaultValue": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Whether to use batching for processing or
streaming. The default is false, which uses streaming. [...]
"batchingIntervalMs": { "index": 15, "kind": "parameter", "displayName":
"Batching Interval Ms", "group": "consumer", "label": "consumer", "required":
false, "type": "integer", "javaType": "java.lang.Integer", "deprecated": false,
"autowired": false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "In consumer batching mode, then this option is
specifying a time in millis, to trigger b [...]
"breakOnFirstError": { "index": 16, "kind": "parameter", "displayName":
"Break On First Error", "group": "consumer", "label": "consumer", "required":
false, "type": "boolean", "javaType": "boolean", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": false,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "This options controls
what happens when a consumer is processing an exchange [...]
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/kafka-component.adoc
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/kafka-component.adoc
index 8519f6a16845..3c7b32d054b6 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/kafka-component.adoc
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/kafka-component.adoc
@@ -344,6 +344,49 @@ exception with the message `KafkaConsumer is not safe for
multi-threaded access`
*Note 2: this is mostly useful with aggregation's completion timeout
strategies.
+=== Where a new consumer group starts reading
+
+The `autoOffsetReset` option controls where a Kafka consumer group begins
reading when it has *no previously committed offset* for the topic partition,
or when the committed offset is no longer available on the broker (e.g. it has
been deleted by retention).
+
+IMPORTANT: `autoOffsetReset` is *only consulted for new or reset consumer
groups*. If the group already has a committed offset, Kafka resumes from that
offset and ignores this option entirely. Setting `by_duration:PT5M` does *not*
rewind an existing consumer group by 5 minutes — it only affects the very first
read of a brand-new group, or a group whose offsets have expired.
+
+The available values are:
+
+`latest` (default):: The consumer starts at the end of the partition (the
latest offset). Messages produced before the consumer first connected are
skipped.
+`earliest`:: The consumer starts at the beginning of the partition. All
retained messages are replayed from the oldest available offset.
+`none`:: An exception is thrown if no previous offset is found. Use this to
detect accidental group-name changes or offset expiry.
+`by_duration:<ISO-8601 duration>`:: *(Kafka 4.0+)* The consumer starts at the
first offset with a timestamp at or after `now - duration`. For example,
`by_duration:PT1H` starts one hour back and `by_duration:P1D` starts one day
back. Like the other values, this only takes effect when the group has no
committed offset.
+
+[tabs]
+====
+Java::
++
+[source,java]
+----
+// New consumer group starting 1 hour back (only on first run)
+from("kafka:orders?brokers=localhost:9092&groupId=reporting&autoOffsetReset=by_duration:PT1H")
+ .log("Order received: ${body}");
+----
+
+YAML::
++
+[source,yaml]
+----
+- route:
+ from:
+ uri: kafka:orders
+ parameters:
+ brokers: "localhost:9092"
+ groupId: reporting
+ autoOffsetReset: "by_duration:PT1H"
+ steps:
+ - log:
+ message: "Order received: ${body}"
+----
+====
+
+NOTE: If you want to reposition an *existing* consumer group on every restart,
use the `seekTo` option (`BEGINNING` or `END`) instead. Unlike
`autoOffsetReset`, `seekTo` applies unconditionally on each consumer start,
regardless of whether committed offsets exist.
+
=== Pausable Consumers
The Kafka component supports pausable consumers. This type of consumer can
pause consuming data based on
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 6317731965ad..c5eda1d3b834 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
@@ -37,7 +37,7 @@
"allowManualCommit": { "index": 10, "kind": "property", "displayName":
"Allow Manual Commit", "group": "consumer", "label": "consumer", "required":
false, "type": "boolean", "javaType": "boolean", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": false,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "Whether to allow doing
manual commits via KafkaManualCommit. If this option is [...]
"autoCommitEnable": { "index": 11, "kind": "property", "displayName":
"Auto Commit Enable", "group": "consumer", "label": "consumer", "required":
false, "type": "boolean", "javaType": "boolean", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": true,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "If true, periodically
commit to ZooKeeper the offset of messages already fetched [...]
"autoCommitIntervalMs": { "index": 12, "kind": "property", "displayName":
"Auto Commit Interval Ms", "group": "consumer", "label": "consumer",
"required": false, "type": "integer", "javaType": "java.lang.Integer",
"deprecated": false, "autowired": false, "secret": false, "defaultValue": 5000,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "The frequency in ms that
the consumer offsets are committed to [...]
- "autoOffsetReset": { "index": 13, "kind": "property", "displayName": "Auto
Offset Reset", "group": "consumer", "label": "consumer", "required": false,
"type": "enum", "javaType": "java.lang.String", "enum": [ "latest", "earliest",
"none" ], "deprecated": false, "autowired": false, "secret": false,
"defaultValue": "latest", "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "What to do when there is no ini [...]
+ "autoOffsetReset": { "index": 13, "kind": "property", "displayName": "Auto
Offset Reset", "group": "consumer", "label": "consumer", "required": false,
"type": "string", "javaType": "java.lang.String", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": "latest",
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "Where a consumer group
starts reading when it has no committed offset, [...]
"batching": { "index": 14, "kind": "property", "displayName": "Batching",
"group": "consumer", "label": "consumer", "required": false, "type": "boolean",
"javaType": "boolean", "deprecated": false, "autowired": false, "secret":
false, "defaultValue": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Whether to use batching for processing or
streaming. The default is false, which uses streaming. I [...]
"batchingIntervalMs": { "index": 15, "kind": "property", "displayName":
"Batching Interval Ms", "group": "consumer", "label": "consumer", "required":
false, "type": "integer", "javaType": "java.lang.Integer", "deprecated": false,
"autowired": false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "In consumer batching mode, then this option is
specifying a time in millis, to trigger ba [...]
"breakOnFirstError": { "index": 16, "kind": "property", "displayName":
"Break On First Error", "group": "consumer", "label": "consumer", "required":
false, "type": "boolean", "javaType": "boolean", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": false,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "This options controls
what happens when a consumer is processing an exchange [...]
@@ -183,7 +183,7 @@
"allowManualCommit": { "index": 10, "kind": "parameter", "displayName":
"Allow Manual Commit", "group": "consumer", "label": "consumer", "required":
false, "type": "boolean", "javaType": "boolean", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": false,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "Whether to allow doing
manual commits via KafkaManualCommit. If this option i [...]
"autoCommitEnable": { "index": 11, "kind": "parameter", "displayName":
"Auto Commit Enable", "group": "consumer", "label": "consumer", "required":
false, "type": "boolean", "javaType": "boolean", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": true,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "If true, periodically
commit to ZooKeeper the offset of messages already fetched [...]
"autoCommitIntervalMs": { "index": 12, "kind": "parameter", "displayName":
"Auto Commit Interval Ms", "group": "consumer", "label": "consumer",
"required": false, "type": "integer", "javaType": "java.lang.Integer",
"deprecated": false, "autowired": false, "secret": false, "defaultValue": 5000,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "The frequency in ms that
the consumer offsets are committed t [...]
- "autoOffsetReset": { "index": 13, "kind": "parameter", "displayName":
"Auto Offset Reset", "group": "consumer", "label": "consumer", "required":
false, "type": "enum", "javaType": "java.lang.String", "enum": [ "latest",
"earliest", "none" ], "deprecated": false, "autowired": false, "secret": false,
"defaultValue": "latest", "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "What to do when there is no in [...]
+ "autoOffsetReset": { "index": 13, "kind": "parameter", "displayName":
"Auto Offset Reset", "group": "consumer", "label": "consumer", "required":
false, "type": "string", "javaType": "java.lang.String", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": "latest",
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "Where a consumer group
starts reading when it has no committed offset, [...]
"batching": { "index": 14, "kind": "parameter", "displayName": "Batching",
"group": "consumer", "label": "consumer", "required": false, "type": "boolean",
"javaType": "boolean", "deprecated": false, "autowired": false, "secret":
false, "defaultValue": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "Whether to use batching for processing or
streaming. The default is false, which uses streaming. [...]
"batchingIntervalMs": { "index": 15, "kind": "parameter", "displayName":
"Batching Interval Ms", "group": "consumer", "label": "consumer", "required":
false, "type": "integer", "javaType": "java.lang.Integer", "deprecated": false,
"autowired": false, "secret": false, "configurationClass":
"org.apache.camel.component.kafka.KafkaConfiguration", "configurationField":
"configuration", "description": "In consumer batching mode, then this option is
specifying a time in millis, to trigger b [...]
"breakOnFirstError": { "index": 16, "kind": "parameter", "displayName":
"Break On First Error", "group": "consumer", "label": "consumer", "required":
false, "type": "boolean", "javaType": "boolean", "deprecated": false,
"autowired": false, "secret": false, "defaultValue": false,
"configurationClass": "org.apache.camel.component.kafka.KafkaConfiguration",
"configurationField": "configuration", "description": "This options controls
what happens when a consumer is processing an exchange [...]
diff --git a/components/camel-kafka/src/main/docs/kafka-component.adoc
b/components/camel-kafka/src/main/docs/kafka-component.adoc
index 8519f6a16845..3c7b32d054b6 100644
--- a/components/camel-kafka/src/main/docs/kafka-component.adoc
+++ b/components/camel-kafka/src/main/docs/kafka-component.adoc
@@ -344,6 +344,49 @@ exception with the message `KafkaConsumer is not safe for
multi-threaded access`
*Note 2: this is mostly useful with aggregation's completion timeout
strategies.
+=== Where a new consumer group starts reading
+
+The `autoOffsetReset` option controls where a Kafka consumer group begins
reading when it has *no previously committed offset* for the topic partition,
or when the committed offset is no longer available on the broker (e.g. it has
been deleted by retention).
+
+IMPORTANT: `autoOffsetReset` is *only consulted for new or reset consumer
groups*. If the group already has a committed offset, Kafka resumes from that
offset and ignores this option entirely. Setting `by_duration:PT5M` does *not*
rewind an existing consumer group by 5 minutes — it only affects the very first
read of a brand-new group, or a group whose offsets have expired.
+
+The available values are:
+
+`latest` (default):: The consumer starts at the end of the partition (the
latest offset). Messages produced before the consumer first connected are
skipped.
+`earliest`:: The consumer starts at the beginning of the partition. All
retained messages are replayed from the oldest available offset.
+`none`:: An exception is thrown if no previous offset is found. Use this to
detect accidental group-name changes or offset expiry.
+`by_duration:<ISO-8601 duration>`:: *(Kafka 4.0+)* The consumer starts at the
first offset with a timestamp at or after `now - duration`. For example,
`by_duration:PT1H` starts one hour back and `by_duration:P1D` starts one day
back. Like the other values, this only takes effect when the group has no
committed offset.
+
+[tabs]
+====
+Java::
++
+[source,java]
+----
+// New consumer group starting 1 hour back (only on first run)
+from("kafka:orders?brokers=localhost:9092&groupId=reporting&autoOffsetReset=by_duration:PT1H")
+ .log("Order received: ${body}");
+----
+
+YAML::
++
+[source,yaml]
+----
+- route:
+ from:
+ uri: kafka:orders
+ parameters:
+ brokers: "localhost:9092"
+ groupId: reporting
+ autoOffsetReset: "by_duration:PT1H"
+ steps:
+ - log:
+ message: "Order received: ${body}"
+----
+====
+
+NOTE: If you want to reposition an *existing* consumer group on every restart,
use the `seekTo` option (`BEGINNING` or `END`) instead. Unlike
`autoOffsetReset`, `seekTo` applies unconditionally on each consumer start,
regardless of whether committed offsets exist.
+
=== Pausable Consumers
The Kafka component supports pausable consumers. This type of consumer can
pause consuming data based on
diff --git
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaConfiguration.java
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaConfiguration.java
index 15d979ef6455..839bc81df83d 100755
---
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaConfiguration.java
+++
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaConfiguration.java
@@ -123,7 +123,7 @@ public class KafkaConfiguration implements Cloneable,
HeaderFilterStrategyAware
@UriParam(label = "consumer", javaType = "java.time.Duration")
private Integer maxPollIntervalMs;
// auto.offset.reset1
- @UriParam(label = "consumer", defaultValue = "latest", enums =
"latest,earliest,none")
+ @UriParam(label = "consumer", defaultValue = "latest")
private String autoOffsetReset = "latest";
// partition.assignment.strategy
@UriParam(label = "consumer", defaultValue =
KafkaConstants.PARTITIONER_RANGE_ASSIGNOR)
@@ -986,9 +986,10 @@ public class KafkaConfiguration implements Cloneable,
HeaderFilterStrategyAware
}
/**
- * What to do when there is no initial offset in ZooKeeper or if an offset
is out of range: earliest : automatically
- * reset the offset to the earliest offset latest: automatically reset the
offset to the latest offset fail: throw
- * exception to the consumer
+ * Where a consumer group starts reading when it has no committed offset,
or the committed offset is out of range.
+ * Valid values are: earliest (seek to the earliest available offset),
latest (seek to the latest offset, the
+ * default), none (throw an exception if no previous offset is found),
by_duration: followed by an ISO-8601 duration
+ * (e.g. by_duration:PT5M or by_duration:P1D; requires Kafka 4.0 or later).
*/
public void setAutoOffsetReset(String autoOffsetReset) {
this.autoOffsetReset = autoOffsetReset;
diff --git
a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaConfigurationTest.java
b/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaConfigurationTest.java
index ae391dc70b7b..e803cae30314 100644
---
a/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaConfigurationTest.java
+++
b/components/camel-kafka/src/test/java/org/apache/camel/component/kafka/KafkaConfigurationTest.java
@@ -16,8 +16,13 @@
*/
package org.apache.camel.component.kafka;
+import java.util.Map;
import java.util.Properties;
+import org.apache.camel.catalog.EndpointValidationResult;
+import org.apache.camel.catalog.RuntimeCamelCatalog;
+import org.apache.camel.catalog.impl.DefaultRuntimeCamelCatalog;
+import org.apache.camel.impl.DefaultCamelContext;
import org.apache.camel.spi.StateRepository;
import org.apache.camel.util.SecurityUtils;
import org.apache.kafka.clients.consumer.ConsumerConfig;
@@ -90,4 +95,16 @@ class KafkaConfigurationTest {
Properties props = config.createConsumerProperties();
assertEquals(131072, props.get(ConsumerConfig.SEND_BUFFER_CONFIG));
}
+
+ @Test
+ void byDurationAutoOffsetResetPassesCatalogValidation() throws Exception {
+ try (DefaultCamelContext context = new DefaultCamelContext()) {
+ RuntimeCamelCatalog catalog = new DefaultRuntimeCamelCatalog();
+ catalog.setCamelContext(context);
+ EndpointValidationResult result
+ = catalog.validateProperties("kafka", Map.of("topic",
"test", "autoOffsetReset", "by_duration:PT5M"));
+ assertTrue(result.isSuccess(),
+ () -> "Expected by_duration:PT5M to pass catalog
validation but got: " + result);
+ }
+ }
}
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 d10ab6998e24..bbcee4eb39c2 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
@@ -306,10 +306,12 @@ public interface KafkaComponentBuilderFactory {
/**
- * What to do when there is no initial offset in ZooKeeper or if an
- * offset is out of range: earliest : automatically reset the offset to
- * the earliest offset latest: automatically reset the offset to the
- * latest offset fail: throw exception to the consumer.
+ * Where a consumer group starts reading when it has no committed
+ * offset, or the committed offset is out of range. Valid values are:
+ * earliest (seek to the earliest available offset), latest (seek to
the
+ * latest offset, the default), none (throw an exception if no previous
+ * offset is found), by_duration: followed by an ISO-8601 duration
(e.g.
+ * by_duration:PT5M or by_duration:P1D; requires Kafka 4.0 or later).
*
* The option is a: <code>java.lang.String</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 afa2f696f9f8..4969eccc4259 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
@@ -455,10 +455,12 @@ public interface KafkaEndpointBuilderFactory {
return this;
}
/**
- * What to do when there is no initial offset in ZooKeeper or if an
- * offset is out of range: earliest : automatically reset the offset to
- * the earliest offset latest: automatically reset the offset to the
- * latest offset fail: throw exception to the consumer.
+ * Where a consumer group starts reading when it has no committed
+ * offset, or the committed offset is out of range. Valid values are:
+ * earliest (seek to the earliest available offset), latest (seek to
the
+ * latest offset, the default), none (throw an exception if no previous
+ * offset is found), by_duration: followed by an ISO-8601 duration
(e.g.
+ * by_duration:PT5M or by_duration:P1D; requires Kafka 4.0 or later).
*
* The option is a: <code>java.lang.String</code> type.
*
diff --git
a/dsl/camel-jbang/camel-jbang-plugin-tui/src/test/java/org/apache/camel/dsl/jbang/core/commands/tui/PropertyCompletionProviderTest.java
b/dsl/camel-jbang/camel-jbang-plugin-tui/src/test/java/org/apache/camel/dsl/jbang/core/commands/tui/PropertyCompletionProviderTest.java
index 73b1851f7639..f27639757ce1 100644
---
a/dsl/camel-jbang/camel-jbang-plugin-tui/src/test/java/org/apache/camel/dsl/jbang/core/commands/tui/PropertyCompletionProviderTest.java
+++
b/dsl/camel-jbang/camel-jbang-plugin-tui/src/test/java/org/apache/camel/dsl/jbang/core/commands/tui/PropertyCompletionProviderTest.java
@@ -253,11 +253,11 @@ class PropertyCompletionProviderTest {
@Test
void componentEnumOptionReturnsValues() {
List<AutocompletePopup.CompletionItem> items
- =
provideValueCompletions("camel.component.kafka.autoOffsetReset");
+ =
provideValueCompletions("camel.component.kafka.compressionCodec");
assertThat(items).isNotEmpty();
- assertThat(items).anyMatch(i -> i.key().equals("latest"));
- assertThat(items).anyMatch(i -> i.key().equals("earliest"));
+ assertThat(items).anyMatch(i -> i.key().equals("none"));
+ assertThat(items).anyMatch(i -> i.key().equals("gzip"));
// each value carries the parent option's description
assertThat(items).allMatch(i -> i.description() != null &&
!i.description().isEmpty());
}
@@ -271,8 +271,10 @@ class PropertyCompletionProviderTest {
@Test
void stringOptionReturnsEmptyValueCompletions() {
// camel.main.name is a string option with no enums
- List<AutocompletePopup.CompletionItem> items =
provideValueCompletions("camel.main.name");
- assertThat(items).isEmpty();
+ assertThat(provideValueCompletions("camel.main.name")).isEmpty();
+ // autoOffsetReset is intentionally a free-form string (not an enum)
because Kafka 4.0
+ // introduced by_duration:<ISO-8601> which cannot be expressed as a
single fixed enum value
+
assertThat(provideValueCompletions("camel.component.kafka.autoOffsetReset")).isEmpty();
}
// --- Helper methods that mirror SourceTab's provider logic ---
diff --git
a/dsl/camel-jbang/camel-jbang-plugin-tui/src/test/java/org/apache/camel/dsl/jbang/core/commands/tui/YamlCompletionTest.java
b/dsl/camel-jbang/camel-jbang-plugin-tui/src/test/java/org/apache/camel/dsl/jbang/core/commands/tui/YamlCompletionTest.java
index fa636d202aab..4135d314e112 100644
---
a/dsl/camel-jbang/camel-jbang-plugin-tui/src/test/java/org/apache/camel/dsl/jbang/core/commands/tui/YamlCompletionTest.java
+++
b/dsl/camel-jbang/camel-jbang-plugin-tui/src/test/java/org/apache/camel/dsl/jbang/core/commands/tui/YamlCompletionTest.java
@@ -332,11 +332,18 @@ class YamlCompletionTest {
@Test
void valueCompletionReturnsEnumValues() {
- List<AutocompletePopup.CompletionItem> items =
provideValueCompletions("kafka", "autoOffsetReset");
+ List<AutocompletePopup.CompletionItem> items =
provideValueCompletions("kafka", "compressionCodec");
+
+ assertThat(items).anyMatch(i -> i.key().equals("none"));
+ assertThat(items).anyMatch(i -> i.key().equals("gzip"));
+ }
- // kafka autoOffsetReset has enum values: latest, earliest, none
- assertThat(items).anyMatch(i -> i.key().equals("latest"));
- assertThat(items).anyMatch(i -> i.key().equals("earliest"));
+ @Test
+ void autoOffsetResetIsStringOptionWithNoEnumCompletions() {
+ // autoOffsetReset is intentionally a free-form string (no fixed enum)
because Kafka 4.0
+ // introduced by_duration:<ISO-8601 duration> which cannot be
expressed as a fixed enum value
+ List<AutocompletePopup.CompletionItem> items =
provideValueCompletions("kafka", "autoOffsetReset");
+ assertThat(items).isEmpty();
}
// --- Property placeholder loading ---