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"