This is an automated email from the ASF dual-hosted git repository.
oscerd pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel-kamelets.git
The following commit(s) were added to refs/heads/main by this push:
new c8cb745fd Fix #2871: surface Kafka record metadata as non-internal
headers (#3032)
c8cb745fd is described below
commit c8cb745fd152e458af0a769a301fbc26ca86acb3
Author: Andrea Cosentino <[email protected]>
AuthorDate: Wed Sep 16 10:19:45 2026 +0200
Fix #2871: surface Kafka record metadata as non-internal headers (#3032)
kafka-source passed the record metadata downstream only as the CamelKafka*
headers, so a consumer wanting the topic, key, partition, offset or
timestamp
had to read Camel internals. This adds a non-internal name for each, with a
CloudEvents counterpart, following infinispan-source and ftp-source:
kafka-topic ce-kafkatopic from CamelKafkaTopic
kafka-key ce-kafkakey from CamelKafkaKey
kafka-partition ce-kafkapartition from CamelKafkaPartition
kafka-offset ce-kafkaoffset from CamelKafkaOffset
kafka-timestamp ce-kafkatimestamp from CamelKafkaTimestamp
The names come from the camel-kafka catalog rather than from assumption, and
only consumer group headers are surfaced; CamelKafkaHeaders, the manual
commit
object and the poll markers are flow control rather than record metadata.
Purely additive. The CamelKafka* headers stay on the exchange, so consumers
already reading them are unaffected. An absent header, CamelKafkaKey for a
record produced without a key, yields an empty header rather than failing
the
exchange.
The kafka-source itest now asserts the new headers arrive at the HTTP sink,
exact for the topic and its ce- form, notEmpty for the rest.
This is the one item of #2871 still outstanding after #2868; the sink side
and
timestamp-router-action were already addressed there.
Co-authored-by: Claude Opus 5 <[email protected]>
---
.../ROOT/partials/kafka-source-description.adoc | 16 +++++++++-
kamelets/kafka-source.kamelet.yaml | 35 ++++++++++++++++++++++
.../kafka/kafka-source-pipe.citrus.it.yaml | 14 +++++++++
3 files changed, 64 insertions(+), 1 deletion(-)
diff --git a/docs/modules/ROOT/partials/kafka-source-description.adoc
b/docs/modules/ROOT/partials/kafka-source-description.adoc
index b17ecf5bb..9c335132d 100644
--- a/docs/modules/ROOT/partials/kafka-source-description.adoc
+++ b/docs/modules/ROOT/partials/kafka-source-description.adoc
@@ -44,7 +44,21 @@ Only `topic` and `bootstrapServers` are required. The
Kamelet supports:
=== Output Format
-The Kamelet outputs Kafka message content and includes Kafka headers and
metadata such as topic, partition, offset, and timestamp.
+The body is the Kafka record value. The record metadata is surfaced as headers
under names
+that are not Camel internals, so a downstream consumer does not have to read
the
+`CamelKafka*` headers directly. Each has a `ce-` prefixed CloudEvents
counterpart as well.
+
+- `kafka-topic` from `CamelKafkaTopic` - the topic the record was consumed
from, which is
+ worth having when the Kamelet subscribes to several topics or to a pattern.
+- `kafka-key` from `CamelKafkaKey` - the record key, absent for records
produced without one.
+- `kafka-partition` from `CamelKafkaPartition` - the partition the record came
from.
+- `kafka-offset` from `CamelKafkaOffset` - the offset within that partition.
+- `kafka-timestamp` from `CamelKafkaTimestamp` - the record timestamp, in
milliseconds since
+ the epoch.
+
+The `CamelKafka*` headers are left on the exchange as well, so consumers
already reading them
+keep working. The Kafka record headers are passed through separately,
controlled by the
+`deserializeHeaders` property.
=== Usage Example
diff --git a/kamelets/kafka-source.kamelet.yaml
b/kamelets/kafka-source.kamelet.yaml
index 234f9469c..1b638a852 100644
--- a/kamelets/kafka-source.kamelet.yaml
+++ b/kamelets/kafka-source.kamelet.yaml
@@ -181,4 +181,39 @@ spec:
steps:
- process:
ref: "{{kafkaHeaderDeserializer}}"
+ # Kafka record metadata, surfaced under names that are not Camel
+ # internals, so a downstream consumer does not have to depend on the
+ # CamelKafka* headers. Plain form plus a CloudEvents one, as
+ # infinispan-source and ftp-source do. The CamelKafka* headers are left
+ # in place, so existing consumers are unaffected.
+ - setHeader:
+ name: kafka-topic
+ simple: "${header[CamelKafkaTopic]}"
+ - setHeader:
+ name: ce-kafkatopic
+ simple: "${header[CamelKafkaTopic]}"
+ - setHeader:
+ name: kafka-key
+ simple: "${header[CamelKafkaKey]}"
+ - setHeader:
+ name: ce-kafkakey
+ simple: "${header[CamelKafkaKey]}"
+ - setHeader:
+ name: kafka-partition
+ simple: "${header[CamelKafkaPartition]}"
+ - setHeader:
+ name: ce-kafkapartition
+ simple: "${header[CamelKafkaPartition]}"
+ - setHeader:
+ name: kafka-offset
+ simple: "${header[CamelKafkaOffset]}"
+ - setHeader:
+ name: ce-kafkaoffset
+ simple: "${header[CamelKafkaOffset]}"
+ - setHeader:
+ name: kafka-timestamp
+ simple: "${header[CamelKafkaTimestamp]}"
+ - setHeader:
+ name: ce-kafkatimestamp
+ simple: "${header[CamelKafkaTimestamp]}"
- to: "kamelet:sink"
diff --git
a/tests/camel-kamelets-itest/src/test/resources/kafka/kafka-source-pipe.citrus.it.yaml
b/tests/camel-kamelets-itest/src/test/resources/kafka/kafka-source-pipe.citrus.it.yaml
index 5f392f370..e01ea5e30 100644
---
a/tests/camel-kamelets-itest/src/test/resources/kafka/kafka-source-pipe.citrus.it.yaml
+++
b/tests/camel-kamelets-itest/src/test/resources/kafka/kafka-source-pipe.citrus.it.yaml
@@ -82,6 +82,20 @@ actions:
headers:
- name: event-source
value: "${kafka.source}"
+ # Record metadata the Kamelet surfaces under non Camel internal
+ # names, so a consumer does not have to read CamelKafka* itself.
+ - name: kafka-topic
+ value: "${kafka.topic}"
+ - name: ce-kafkatopic
+ value: "${kafka.topic}"
+ - name: kafka-key
+ value: "@notEmpty()@"
+ - name: kafka-partition
+ value: "@notEmpty()@"
+ - name: kafka-offset
+ value: "@notEmpty()@"
+ - name: kafka-timestamp
+ value: "@notEmpty()@"
body:
data: |
${kafka.message}