From b4b2dcfc1167767e9f689169fcb08fb2922bfd9c Mon Sep 17 00:00:00 2001 From: Salvatore Mongiardo Date: Tue, 29 Sep 2026 16:40:40 +0200 Subject: [PATCH] CAMEL-25091: camel-kafka: expose by_duration as valid autoOffsetReset value (Kafka 4.0+) Kafka 4.0 introduced the by_duration: auto offset reset strategy (KAFKA-18013), which seeks the consumer to the offset at (now - duration) when it starts. For example, by_duration:PT5M positions the consumer at the offset from 5 minutes before startup. Camel already passes autoOffsetReset through to ConsumerConfig without runtime validation, so by_duration:PT5M works on any build backed by kafka-clients 4.x (available since Camel 4.19 via CAMEL-23086). The only gap was documentation and tooling visibility. Changes: - Remove the enums attribute from @UriParam on autoOffsetReset so that by_duration: values pass catalog and camel validate checks (a bare enum entry "by_duration" would be invalid; the real form is parameterised and cannot be expressed as a fixed enum value) - Rewrite the autoOffsetReset description: drop the stale ZooKeeper reference, correct "fail" to "none", and add by_duration documentation - Regenerate the catalog mirror (camel-catalog/components/kafka.json) and update the endpoint-dsl and component-dsl factory javadocs - Add a catalog validation test using DefaultRuntimeCamelCatalog (already on classpath via camel-core, no new dependency needed) that by_duration:PT5M is accepted by validateProperties, verifying the end-to-end user-facing behaviour --- .../camel/catalog/components/kafka.json | 4 +- .../camel/catalog/docs/kafka-component.adoc | 43 +++++++++++++++++++ .../apache/camel/component/kafka/kafka.json | 4 +- .../src/main/docs/kafka-component.adoc | 43 +++++++++++++++++++ .../component/kafka/KafkaConfiguration.java | 9 ++-- .../kafka/KafkaConfigurationTest.java | 17 ++++++++ .../dsl/KafkaComponentBuilderFactory.java | 10 +++-- .../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 94c980effce0f..416a4ed1765d7 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 enabled then an instance of KafkaManualCommit is stored on the Exchange message header, which allows end users to access this API and perform manual offset commits via the Kafka consumer." }, "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 by the consumer. This committed offset will be used when the process fails as the position from which the new consumer will begin." }, "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 zookeeper." }, - "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 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" }, + "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, 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)." }, "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. In streaming mode, then a single kafka record is processed per Camel exchange in the message body. In batching mode, then Camel groups many kafka records together as a List objects in the message body. The option maxPollRecords is used to define the number of records to group together in batching mode." }, "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 batch completion eager when the current batch size has not reached the maximum size defined by maxPollRecords. Notice the trigger is not exact at the given interval, as this can only happen between kafka polls (see pollTimeoutMs option). So for example setting this to 10000, then the trigger happens in the interval 10000 pollTimeoutMs. The default value for pollTimeoutMs is 5000, so this would mean a trigger interval at about every 15 seconds." }, "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 and it fails. If the option is false then the consumer continues to the next message and processes it. If the option is true then the consumer breaks out. Using the default NoopCommitManager will cause the consumer to not commit the offset so that the message is re-attempted. The consumer should use the KafkaManualCommit to determine the best way to handle the message. Using either the SyncCommitManager or the AsyncCommitManager, the consumer will seek back to the offset of the message that caused a failure, and then re-attempt to process this message. However, this can lead to endless processing of the same message if it's bound to fail every time, e.g., a poison message. Therefore, it's recommended to deal with that, for example, by using Camel's error handler." }, @@ -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 is enabled then an instance of KafkaManualCommit is stored on the Exchange message header, which allows end users to access this API and perform manual offset commits via the Kafka consumer." }, "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 by the consumer. This committed offset will be used when the process fails as the position from which the new consumer will begin." }, "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 to zookeeper." }, - "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 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" }, + "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, 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)." }, "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. In streaming mode, then a single kafka record is processed per Camel exchange in the message body. In batching mode, then Camel groups many kafka records together as a List objects in the message body. The option maxPollRecords is used to define the number of records to group together in batching mode." }, "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 batch completion eager when the current batch size has not reached the maximum size defined by maxPollRecords. Notice the trigger is not exact at the given interval, as this can only happen between kafka polls (see pollTimeoutMs option). So for example setting this to 10000, then the trigger happens in the interval 10000 pollTimeoutMs. The default value for pollTimeoutMs is 5000, so this would mean a trigger interval at about every 15 seconds." }, "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 and it fails. If the option is false then the consumer continues to the next message and processes it. If the option is true then the consumer breaks out. Using the default NoopCommitManager will cause the consumer to not commit the offset so that the message is re-attempted. The consumer should use the KafkaManualCommit to determine the best way to handle the message. Using either the SyncCommitManager or the AsyncCommitManager, the consumer will seek back to the offset of the message that caused a failure, and then re-attempt to process this message. However, this can lead to endless processing of the same message if it's bound to fail every time, e.g., a poison message. Therefore, it's recommended to deal with that, for example, by using Camel's error handler." }, 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 9a9b669a78ec3..3d1cb5eb2d8cf 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:`:: *(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 94c980effce0f..416a4ed1765d7 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 enabled then an instance of KafkaManualCommit is stored on the Exchange message header, which allows end users to access this API and perform manual offset commits via the Kafka consumer." }, "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 by the consumer. This committed offset will be used when the process fails as the position from which the new consumer will begin." }, "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 zookeeper." }, - "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 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" }, + "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, 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)." }, "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. In streaming mode, then a single kafka record is processed per Camel exchange in the message body. In batching mode, then Camel groups many kafka records together as a List objects in the message body. The option maxPollRecords is used to define the number of records to group together in batching mode." }, "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 batch completion eager when the current batch size has not reached the maximum size defined by maxPollRecords. Notice the trigger is not exact at the given interval, as this can only happen between kafka polls (see pollTimeoutMs option). So for example setting this to 10000, then the trigger happens in the interval 10000 pollTimeoutMs. The default value for pollTimeoutMs is 5000, so this would mean a trigger interval at about every 15 seconds." }, "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 and it fails. If the option is false then the consumer continues to the next message and processes it. If the option is true then the consumer breaks out. Using the default NoopCommitManager will cause the consumer to not commit the offset so that the message is re-attempted. The consumer should use the KafkaManualCommit to determine the best way to handle the message. Using either the SyncCommitManager or the AsyncCommitManager, the consumer will seek back to the offset of the message that caused a failure, and then re-attempt to process this message. However, this can lead to endless processing of the same message if it's bound to fail every time, e.g., a poison message. Therefore, it's recommended to deal with that, for example, by using Camel's error handler." }, @@ -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 is enabled then an instance of KafkaManualCommit is stored on the Exchange message header, which allows end users to access this API and perform manual offset commits via the Kafka consumer." }, "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 by the consumer. This committed offset will be used when the process fails as the position from which the new consumer will begin." }, "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 to zookeeper." }, - "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 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" }, + "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, 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)." }, "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. In streaming mode, then a single kafka record is processed per Camel exchange in the message body. In batching mode, then Camel groups many kafka records together as a List objects in the message body. The option maxPollRecords is used to define the number of records to group together in batching mode." }, "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 batch completion eager when the current batch size has not reached the maximum size defined by maxPollRecords. Notice the trigger is not exact at the given interval, as this can only happen between kafka polls (see pollTimeoutMs option). So for example setting this to 10000, then the trigger happens in the interval 10000 pollTimeoutMs. The default value for pollTimeoutMs is 5000, so this would mean a trigger interval at about every 15 seconds." }, "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 and it fails. If the option is false then the consumer continues to the next message and processes it. If the option is true then the consumer breaks out. Using the default NoopCommitManager will cause the consumer to not commit the offset so that the message is re-attempted. The consumer should use the KafkaManualCommit to determine the best way to handle the message. Using either the SyncCommitManager or the AsyncCommitManager, the consumer will seek back to the offset of the message that caused a failure, and then re-attempt to process this message. However, this can lead to endless processing of the same message if it's bound to fail every time, e.g., a poison message. Therefore, it's recommended to deal with that, for example, by using Camel's error handler." }, diff --git a/components/camel-kafka/src/main/docs/kafka-component.adoc b/components/camel-kafka/src/main/docs/kafka-component.adoc index 9a9b669a78ec3..3d1cb5eb2d8cf 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:`:: *(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 4050826684207..230a76ca5ee5a 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 @@ -124,7 +124,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) @@ -1019,9 +1019,10 @@ public String getAutoOffsetReset() { } /** - * 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 ae391dc70b7b4..e803cae303144 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 @@ void sendBufferBytesAppliedToConsumerWithoutSsl() { 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 d10ab6998e243..bbcee4eb39c24 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 @@ default KafkaComponentBuilder autoCommitIntervalMs(java.lang.Integer autoCommitI /** - * 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 afa2f696f9f82..4969eccc42595 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 @@ default KafkaEndpointConsumerBuilder autoCommitIntervalMs(String autoCommitInter 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: java.lang.String 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 73b1851f7639e..f27639757ce12 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 @@ void enumOptionReturnsEnumValues() { @Test void componentEnumOptionReturnsValues() { List 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 @@ void unknownKeyReturnsEmptyValueCompletions() { @Test void stringOptionReturnsEmptyValueCompletions() { // camel.main.name is a string option with no enums - List 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: 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 79b51dfdeac20..476e5f0efc4ab 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 @@ -418,11 +418,18 @@ void valueCompletionReturnsBooleanValues() { @Test void valueCompletionReturnsEnumValues() { - List items = provideValueCompletions("kafka", "autoOffsetReset"); + List 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: which cannot be expressed as a fixed enum value + List items = provideValueCompletions("kafka", "autoOffsetReset"); + assertThat(items).isEmpty(); } // --- Property placeholder loading ---