From 7259e5eb5e21d629a74c88d5cf201757fcfa0e84 Mon Sep 17 00:00:00 2001 From: Nathan Xu Date: Wed, 13 Dec 2023 10:35:57 -0500 Subject: [PATCH] Reuse Header in AggregatingReplyingKafkaTemplate First of all we extract `Header correlation = record.headers().lastHeader(correlationHeaderName);`. Then we convert its value back for a new header we add back to the aggregated records in the list. Therefore no need in extra object and conversion logic. --- .../requestreply/AggregatingReplyingKafkaTemplate.java | 7 +------ 1 file changed, 1 insertion(+), 6 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/requestreply/AggregatingReplyingKafkaTemplate.java b/spring-kafka/src/main/java/org/springframework/kafka/requestreply/AggregatingReplyingKafkaTemplate.java index 0bfebbe0..ba63cabc 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/requestreply/AggregatingReplyingKafkaTemplate.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/requestreply/AggregatingReplyingKafkaTemplate.java @@ -34,7 +34,6 @@ import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.header.Header; -import org.apache.kafka.common.header.internals.RecordHeader; import org.springframework.kafka.core.ProducerFactory; import org.springframework.kafka.listener.BatchConsumerAwareMessageListener; @@ -148,11 +147,7 @@ public class AggregatingReplyingKafkaTemplate if (this.releaseStrategy.test(list, false)) { ConsumerRecord>> done = new ConsumerRecord<>(AGGREGATED_RESULTS_TOPIC, 0, 0L, null, list); - done.headers() - .add(new RecordHeader(correlationHeaderName, - isBinaryCorrelation() - ? ((CorrelationKey) correlationId).getCorrelationId() - : ((String) correlationId).getBytes(StandardCharsets.UTF_8))); + done.headers().add(correlation); this.pending.remove(correlationId); checkOffsetsAndCommitIfNecessary(list, consumer); completed.add(done);