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());