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 2e1944706 Fix #2871: surface Kafka record metadata on the remaining 
single record sources (#3036)
2e1944706 is described below

commit 2e19447069ea4c1593669e91186cae723ce31333
Author: Andrea Cosentino <[email protected]>
AuthorDate: Wed Sep 16 10:36:13 2026 +0200

    Fix #2871: surface Kafka record metadata on the remaining single record 
sources (#3036)
    
    Follow-up to #3032, which added the mapping to kafka-source only. The other
    single record Kafka sources consume the same per-record CamelKafka* headers 
and
    had the same gap:
    
      kafka-apicurio-registry-not-secured-source
      kafka-azure-schema-registry-source
      kafka-not-secured-apicurio-registry-source
      kafka-not-secured-apicurio-registry-json-source
    
    Each gets the same block as kafka-source, mapping CamelKafkaTopic, Key,
    Partition, Offset and Timestamp to kafka-topic, kafka-key, kafka-partition,
    kafka-offset and kafka-timestamp, plus a ce- CloudEvents counterpart each.
    
    The four kafka-batch-* sources are deliberately left out. They run the 
consumer
    with batching enabled, so the body is a list of records and there is no 
single
    record whose topic, partition or offset the headers could describe; mapping 
them
    there would be misleading rather than useful.
    
    Purely additive, as in #3032: the CamelKafka* headers stay on the exchange.
    
    These four have no Citrus tests, since they need a live Apicurio or Azure 
Schema
    Registry, so the mapping was verified by loading all four templates in one 
Camel
    context and driving a record through each. The behaviour itself is covered 
end
    to end against a real broker by the kafka-source-pipe-test strengthened in 
#3032.
    
    Co-authored-by: Claude Opus 5 <[email protected]>
---
 ...io-registry-not-secured-source-description.adoc | 14 +++++++++
 ...a-azure-schema-registry-source-description.adoc | 14 +++++++++
 ...-apicurio-registry-json-source-description.adoc | 14 +++++++++
 ...cured-apicurio-registry-source-description.adoc | 14 +++++++++
 ...icurio-registry-not-secured-source.kamelet.yaml | 35 ++++++++++++++++++++++
 ...kafka-azure-schema-registry-source.kamelet.yaml | 35 ++++++++++++++++++++++
 ...ured-apicurio-registry-json-source.kamelet.yaml | 35 ++++++++++++++++++++++
 ...t-secured-apicurio-registry-source.kamelet.yaml | 35 ++++++++++++++++++++++
 8 files changed, 196 insertions(+)

diff --git 
a/docs/modules/ROOT/partials/kafka-apicurio-registry-not-secured-source-description.adoc
 
b/docs/modules/ROOT/partials/kafka-apicurio-registry-not-secured-source-description.adoc
index a86fadec9..0172ba9e6 100644
--- 
a/docs/modules/ROOT/partials/kafka-apicurio-registry-not-secured-source-description.adoc
+++ 
b/docs/modules/ROOT/partials/kafka-apicurio-registry-not-secured-source-description.adoc
@@ -12,6 +12,20 @@ This Kamelet connects to Kafka using:
 
 The Kamelet consumes messages from Kafka topics with schema validation through 
Apicurio Registry and produces the message data in the configured format.
 
+The Kafka 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.
+- `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.
+
 === Configuration
 
 The Kamelet requires Kafka and Apicurio Registry connection parameters:
diff --git 
a/docs/modules/ROOT/partials/kafka-azure-schema-registry-source-description.adoc
 
b/docs/modules/ROOT/partials/kafka-azure-schema-registry-source-description.adoc
index 14aff7a11..25a60a9df 100644
--- 
a/docs/modules/ROOT/partials/kafka-azure-schema-registry-source-description.adoc
+++ 
b/docs/modules/ROOT/partials/kafka-azure-schema-registry-source-description.adoc
@@ -12,6 +12,20 @@ This Kamelet connects to Kafka using:
 
 The Kamelet consumes messages from Kafka topics with schema validation through 
Azure Schema Registry and produces the message data in the configured format.
 
+The Kafka 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.
+- `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.
+
 === Configuration
 
 The Kamelet requires Kafka and Azure Schema Registry connection parameters:
diff --git 
a/docs/modules/ROOT/partials/kafka-not-secured-apicurio-registry-json-source-description.adoc
 
b/docs/modules/ROOT/partials/kafka-not-secured-apicurio-registry-json-source-description.adoc
index 2604288fe..01ff99df6 100644
--- 
a/docs/modules/ROOT/partials/kafka-not-secured-apicurio-registry-json-source-description.adoc
+++ 
b/docs/modules/ROOT/partials/kafka-not-secured-apicurio-registry-json-source-description.adoc
@@ -12,6 +12,20 @@ This Kamelet connects to Kafka using appropriate security 
mechanisms based on th
 
 The Kamelet consumes messages from Kafka topics and produces the message data 
in the configured format.
 
+The Kafka 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.
+- `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.
+
 === Configuration
 
 The Kamelet requires Kafka connection parameters:
diff --git 
a/docs/modules/ROOT/partials/kafka-not-secured-apicurio-registry-source-description.adoc
 
b/docs/modules/ROOT/partials/kafka-not-secured-apicurio-registry-source-description.adoc
index 15ef6b4aa..b595fdcdb 100644
--- 
a/docs/modules/ROOT/partials/kafka-not-secured-apicurio-registry-source-description.adoc
+++ 
b/docs/modules/ROOT/partials/kafka-not-secured-apicurio-registry-source-description.adoc
@@ -12,6 +12,20 @@ This Kamelet connects to Kafka using appropriate security 
mechanisms based on th
 
 The Kamelet consumes messages from Kafka topics and produces the message data 
in the configured format.
 
+The Kafka 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.
+- `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.
+
 === Configuration
 
 The Kamelet requires Kafka connection parameters:
diff --git a/kamelets/kafka-apicurio-registry-not-secured-source.kamelet.yaml 
b/kamelets/kafka-apicurio-registry-not-secured-source.kamelet.yaml
index 44247ff0a..ddf44a23a 100644
--- a/kamelets/kafka-apicurio-registry-not-secured-source.kamelet.yaml
+++ b/kamelets/kafka-apicurio-registry-not-secured-source.kamelet.yaml
@@ -115,4 +115,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/kamelets/kafka-azure-schema-registry-source.kamelet.yaml 
b/kamelets/kafka-azure-schema-registry-source.kamelet.yaml
index 8e8344560..900ee9216 100644
--- a/kamelets/kafka-azure-schema-registry-source.kamelet.yaml
+++ b/kamelets/kafka-azure-schema-registry-source.kamelet.yaml
@@ -150,4 +150,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/kamelets/kafka-not-secured-apicurio-registry-json-source.kamelet.yaml 
b/kamelets/kafka-not-secured-apicurio-registry-json-source.kamelet.yaml
index 4c5a6e718..fa95219f4 100644
--- a/kamelets/kafka-not-secured-apicurio-registry-json-source.kamelet.yaml
+++ b/kamelets/kafka-not-secured-apicurio-registry-json-source.kamelet.yaml
@@ -158,4 +158,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/kamelets/kafka-not-secured-apicurio-registry-source.kamelet.yaml 
b/kamelets/kafka-not-secured-apicurio-registry-source.kamelet.yaml
index 4eeb0c89d..75ca2e508 100644
--- a/kamelets/kafka-not-secured-apicurio-registry-source.kamelet.yaml
+++ b/kamelets/kafka-not-secured-apicurio-registry-source.kamelet.yaml
@@ -167,4 +167,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"

Reply via email to