This is an automated email from the ASF dual-hosted git repository.

davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/main by this push:
     new ed8a9624f1bc CAMEL-24779: camel-kafka - reduce per-message allocations 
on hot paths
ed8a9624f1bc is described below

commit ed8a9624f1bcf26931b60b0d4c93838032078fc0
Author: Andrea Cosentino <[email protected]>
AuthorDate: Thu Sep 17 08:03:43 2026 +0200

    CAMEL-24779: camel-kafka - reduce per-message allocations on hot paths
    
    Three behaviour-preserving clean-ups that remove per-message allocations
    in camel-kafka. No public API or option changes.
    
    Producer: the single-message async path in KafkaProducer.process passed
    a non-null key to doSend, which allocated a KafkaProducerMetadataCallBack
    and a DelegatingCallback per message even though the parent
    KafkaProducerCallBack already records metadata and exceptions on the same
    exchange. It now sends with the parent callback alone; the batch path is
    unchanged.
    
    Consumer: KafkaRecordProcessor.propagateHeaders built a Stream, spliterator
    and two capturing lambdas for every record and re-resolved exchange.getIn()
    per header. It now uses a plain loop with a hoisted Message.
    
    Transforms: HoistField, MaskField, ExtractField, ReplaceField,
    MessageTimestampRouter and ValueToKey constructed a new ObjectMapper on
    every invocation. They now share a single static instance.
    
    Closes #26522
    
    Co-authored-by: Claude <[email protected]>
---
 .../org/apache/camel/component/kafka/KafkaProducer.java    |  5 ++++-
 .../kafka/consumer/support/KafkaRecordProcessor.java       | 14 ++++++++------
 .../camel/component/kafka/transform/ExtractField.java      |  5 +++--
 .../apache/camel/component/kafka/transform/HoistField.java |  5 +++--
 .../apache/camel/component/kafka/transform/MaskField.java  |  8 ++++----
 .../component/kafka/transform/MessageTimestampRouter.java  |  5 +++--
 .../camel/component/kafka/transform/ReplaceField.java      |  9 +++++----
 .../apache/camel/component/kafka/transform/ValueToKey.java |  5 +++--
 8 files changed, 33 insertions(+), 23 deletions(-)

diff --git 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java
 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java
index 7334c4b93756..cb6b92701643 100755
--- 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java
+++ 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/KafkaProducer.java
@@ -464,7 +464,10 @@ public class KafkaProducer extends DefaultAsyncProducer 
implements RouteIdAware
                 processIterableAsync(exchange, producerCallBack, message);
             } else {
                 final ProducerRecord<Object, Object> record = 
createRecord(exchange, message);
-                doSend(exchange, record, producerCallBack);
+                // Single message: the parent KafkaProducerCallBack already 
records the metadata and any
+                // exception on this exchange, so pass a null key to skip the 
redundant per-record metadata
+                // callback (avoids two short-lived allocations per message) 
(CAMEL-24779).
+                doSend(null, record, producerCallBack);
             }
 
             return producerCallBack.allSent();
diff --git 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/support/KafkaRecordProcessor.java
 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/support/KafkaRecordProcessor.java
index dd330042e463..35f8c45d640b 100644
--- 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/support/KafkaRecordProcessor.java
+++ 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/consumer/support/KafkaRecordProcessor.java
@@ -17,8 +17,6 @@
 
 package org.apache.camel.component.kafka.consumer.support;
 
-import java.util.stream.StreamSupport;
-
 import org.apache.camel.Exchange;
 import org.apache.camel.Message;
 import org.apache.camel.component.kafka.KafkaConfiguration;
@@ -60,10 +58,14 @@ public abstract class KafkaRecordProcessor {
 
         HeaderFilterStrategy headerFilterStrategy = 
configuration.getHeaderFilterStrategy();
         KafkaHeaderDeserializer headerDeserializer = 
configuration.getHeaderDeserializer();
+        Message in = exchange.getIn();
 
-        StreamSupport.stream(consumerRecord.headers().spliterator(), false)
-                .filter(header -> shouldBeFiltered(header, exchange, 
headerFilterStrategy))
-                .forEach(header -> exchange.getIn().setHeader(header.key(),
-                        headerDeserializer.deserialize(header.key(), 
header.value())));
+        // Iterate the record headers directly instead of allocating a Stream, 
spliterator and lambdas per
+        // consumed record; getIn() is resolved once rather than for every 
header (CAMEL-24779).
+        for (Header header : consumerRecord.headers()) {
+            if (shouldBeFiltered(header, exchange, headerFilterStrategy)) {
+                in.setHeader(header.key(), 
headerDeserializer.deserialize(header.key(), header.value()));
+            }
+        }
     }
 }
diff --git 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/ExtractField.java
 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/ExtractField.java
index 430c25880315..24ddcb96406a 100644
--- 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/ExtractField.java
+++ 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/ExtractField.java
@@ -27,6 +27,8 @@ import org.apache.camel.Processor;
 
 public class ExtractField implements Processor {
 
+    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
+
     String field;
     String headerOutputName;
     boolean headerOutput;
@@ -52,7 +54,6 @@ public class ExtractField implements Processor {
 
     @Override
     public void process(Exchange ex) throws InvalidPayloadException {
-        ObjectMapper mapper = new ObjectMapper();
         JsonNode jsonNodeBody = ex.getMessage().getBody(JsonNode.class);
 
         if (jsonNodeBody == null) {
@@ -60,7 +61,7 @@ public class ExtractField implements Processor {
 
         }
 
-        Map<Object, Object> body = mapper.convertValue(jsonNodeBody, new 
TypeReference<Map<Object, Object>>() {
+        Map<Object, Object> body = OBJECT_MAPPER.convertValue(jsonNodeBody, 
new TypeReference<Map<Object, Object>>() {
         });
         if (!headerOutput || (strictHeaderCheck && checkHeaderExistence(ex))) {
             ex.getMessage().setBody(body.get(field));
diff --git 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/HoistField.java
 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/HoistField.java
index 4f9d681bfaee..8032f128f944 100644
--- 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/HoistField.java
+++ 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/HoistField.java
@@ -27,12 +27,13 @@ import org.apache.camel.InvalidPayloadException;
 
 public class HoistField {
 
+    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
+
     public JsonNode process(@ExchangeProperty("field") String field, Exchange 
ex) throws InvalidPayloadException {
-        ObjectMapper mapper = new ObjectMapper();
         Object body = ex.getMessage().getBody();
         Map<Object, Object> updatedBody = new HashMap<>();
         updatedBody.put(field, body);
-        return mapper.valueToTree(updatedBody);
+        return OBJECT_MAPPER.valueToTree(updatedBody);
     }
 
 }
diff --git 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/MaskField.java
 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/MaskField.java
index e40f44adf2a8..0ac27acc11a2 100644
--- 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/MaskField.java
+++ 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/MaskField.java
@@ -32,6 +32,7 @@ import org.apache.camel.util.ObjectHelper;
 
 public class MaskField {
 
+    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
     private static final Map<Class<?>, Function<String, ?>> MAPPING_FUNC = new 
HashMap<>();
     private static final Map<Class<?>, Object> BASIC_MAPPING = new HashMap<>();
 
@@ -62,10 +63,9 @@ public class MaskField {
     public JsonNode process(
             @ExchangeProperty("fields") String fields, 
@ExchangeProperty("replacement") String replacement, Exchange ex)
             throws InvalidPayloadException {
-        ObjectMapper mapper = new ObjectMapper();
         List<String> splittedFields = new ArrayList<>();
         JsonNode jsonNodeBody = ex.getMessage().getBody(JsonNode.class);
-        Map<Object, Object> body = mapper.convertValue(jsonNodeBody, new 
TypeReference<Map<Object, Object>>() {
+        Map<Object, Object> body = OBJECT_MAPPER.convertValue(jsonNodeBody, 
new TypeReference<Map<Object, Object>>() {
         });
         if (ObjectHelper.isNotEmpty(fields)) {
             splittedFields = 
Arrays.stream(fields.split(",")).collect(Collectors.toList());
@@ -79,9 +79,9 @@ public class MaskField {
                     filterNames(fieldName, splittedFields) ? 
masked(origFieldValue, replacement) : origFieldValue);
         }
         if (!updatedBody.isEmpty()) {
-            return mapper.valueToTree(updatedBody);
+            return OBJECT_MAPPER.valueToTree(updatedBody);
         } else {
-            return mapper.valueToTree(body);
+            return OBJECT_MAPPER.valueToTree(body);
         }
     }
 
diff --git 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/MessageTimestampRouter.java
 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/MessageTimestampRouter.java
index 6225f938392d..022345e7f429 100644
--- 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/MessageTimestampRouter.java
+++ 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/MessageTimestampRouter.java
@@ -33,6 +33,8 @@ import org.apache.camel.util.ObjectHelper;
 
 public class MessageTimestampRouter {
 
+    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
+
     public void process(
             @ExchangeProperty("topicFormat") String topicFormat, 
@ExchangeProperty("timestampFormat") String timestampFormat,
             @ExchangeProperty("timestampKeys") String timestampKeys,
@@ -45,10 +47,9 @@ public class MessageTimestampRouter {
         final SimpleDateFormat fmt = new SimpleDateFormat(timestampFormat);
         fmt.setTimeZone(TimeZone.getTimeZone("UTC"));
 
-        ObjectMapper mapper = new ObjectMapper();
         List<String> splittedKeys = new ArrayList<>();
         JsonNode jsonNodeBody = ex.getMessage().getBody(JsonNode.class);
-        Map<Object, Object> body = mapper.convertValue(jsonNodeBody, new 
TypeReference<Map<Object, Object>>() {
+        Map<Object, Object> body = OBJECT_MAPPER.convertValue(jsonNodeBody, 
new TypeReference<Map<Object, Object>>() {
         });
         if (ObjectHelper.isNotEmpty(timestampKeys)) {
             splittedKeys = 
Arrays.stream(timestampKeys.split(",")).collect(Collectors.toList());
diff --git 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/ReplaceField.java
 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/ReplaceField.java
index c9d0499373af..71a278a17d93 100644
--- 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/ReplaceField.java
+++ 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/ReplaceField.java
@@ -29,16 +29,17 @@ import org.apache.camel.util.ObjectHelper;
 
 public class ReplaceField {
 
+    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
+
     public JsonNode process(
             @ExchangeProperty("enabled") String enabled, 
@ExchangeProperty("disabled") String disabled,
             @ExchangeProperty("renames") String renames, Exchange ex)
             throws InvalidPayloadException {
-        ObjectMapper mapper = new ObjectMapper();
         List<String> enabledFields = new ArrayList<>();
         List<String> disabledFields = new ArrayList<>();
         List<String> renameFields = new ArrayList<>();
         JsonNode jsonNodeBody = ex.getMessage().getBody(JsonNode.class);
-        Map<Object, Object> body = mapper.convertValue(jsonNodeBody, new 
TypeReference<Map<Object, Object>>() {
+        Map<Object, Object> body = OBJECT_MAPPER.convertValue(jsonNodeBody, 
new TypeReference<Map<Object, Object>>() {
         });
         if (ObjectHelper.isNotEmpty(enabled) && 
!enabled.equalsIgnoreCase("all")) {
             enabledFields = 
Arrays.stream(enabled.split(",")).collect(Collectors.toList());
@@ -63,9 +64,9 @@ public class ReplaceField {
             }
         }
         if (!updatedBody.isEmpty()) {
-            return mapper.valueToTree(updatedBody);
+            return OBJECT_MAPPER.valueToTree(updatedBody);
         } else {
-            return mapper.valueToTree(body);
+            return OBJECT_MAPPER.valueToTree(body);
         }
     }
 
diff --git 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/ValueToKey.java
 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/ValueToKey.java
index b2fc89b9a2c8..c2660a3f2462 100644
--- 
a/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/ValueToKey.java
+++ 
b/components/camel-kafka/src/main/java/org/apache/camel/component/kafka/transform/ValueToKey.java
@@ -30,11 +30,12 @@ import org.apache.camel.util.ObjectHelper;
 
 public class ValueToKey {
 
+    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
+
     public void process(@ExchangeProperty("fields") String fields, Exchange 
ex) throws InvalidPayloadException {
         List<String> splittedFields = new ArrayList<>();
-        ObjectMapper mapper = new ObjectMapper();
         JsonNode jsonNodeBody = ex.getMessage().getBody(JsonNode.class);
-        Map<Object, Object> body = mapper.convertValue(jsonNodeBody, new 
TypeReference<Map<Object, Object>>() {
+        Map<Object, Object> body = OBJECT_MAPPER.convertValue(jsonNodeBody, 
new TypeReference<Map<Object, Object>>() {
         });
         if (ObjectHelper.isNotEmpty(fields)) {
             splittedFields = 
Arrays.stream(fields.split(",")).collect(Collectors.toList());

Reply via email to