From b7de020e165fe07ca2c5b30d2bc2caca26a6f66a Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Thu, 15 Feb 2024 16:25:12 -0500 Subject: [PATCH] Fixing concurrency issue in Kafka Streams binder For more details: https://github.com/spring-cloud/spring-cloud-stream/commit/e7c1c7f7690a790aac4f83469b61bfad60cf9b15 --- .../streams/KafkaStreamsMessageConversionDelegate.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java index 60cf9a5b8..4ffbdfcdb 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsMessageConversionDelegate.java @@ -90,7 +90,7 @@ public class KafkaStreamsMessageConversionDelegate { String contentType = this.kstreamBindingInformationCatalogue .getContentType(outboundBindTarget); MessageConverter messageConverter = this.compositeMessageConverter; - final PerRecordContentTypeHolder perRecordContentTypeHolder = new PerRecordContentTypeHolder(); + final ThreadLocal perRecordContentTypeHolderThreadLocal = ThreadLocal.withInitial(PerRecordContentTypeHolder::new); final KStream kStreamWithEnrichedHeaders = outboundBindTarget .filter((k, v) -> v != null) @@ -103,7 +103,7 @@ public class KafkaStreamsMessageConversionDelegate { } MessageHeaders messageHeaders = new MessageHeaders(headers); final Message convertedMessage = messageConverter.toMessage(message.getPayload(), messageHeaders); - perRecordContentTypeHolder.setContentType((String) messageHeaders.get(MessageHeaders.CONTENT_TYPE)); + perRecordContentTypeHolderThreadLocal.get().setContentType((String) messageHeaders.get(MessageHeaders.CONTENT_TYPE)); return convertedMessage.getPayload(); }); @@ -118,12 +118,12 @@ public class KafkaStreamsMessageConversionDelegate { @Override public void process(Object key, Object value) { - if (perRecordContentTypeHolder.contentType != null) { + if (perRecordContentTypeHolderThreadLocal.get().contentType != null) { this.context.headers().remove(MessageHeaders.CONTENT_TYPE); final Header header; try { header = new RecordHeader(MessageHeaders.CONTENT_TYPE, - new ObjectMapper().writeValueAsBytes(perRecordContentTypeHolder.contentType)); + new ObjectMapper().writeValueAsBytes(perRecordContentTypeHolderThreadLocal.get().contentType)); this.context.headers().add(header); } catch (Exception e) { @@ -131,7 +131,7 @@ public class KafkaStreamsMessageConversionDelegate { LOG.debug("Could not add content type header"); } } - perRecordContentTypeHolder.unsetContentType(); + perRecordContentTypeHolderThreadLocal.get().unsetContentType(); } }