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 509b42025..c4c5bc211 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 @@ -1,5 +1,5 @@ /* - * Copyright 2018-2023 the original author or authors. + * Copyright 2018-2024 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -53,6 +53,7 @@ import org.springframework.util.StringUtils; * * @author Soby Chacko * @author Byungjun You + * @author Lazare Giorgobiani */ public class KafkaStreamsMessageConversionDelegate { @@ -92,7 +93,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) @@ -105,7 +106,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 Objects.requireNonNull(convertedMessage).getPayload(); }); @@ -120,12 +121,12 @@ public class KafkaStreamsMessageConversionDelegate { @Override public void process(Record record) { - if (perRecordContentTypeHolder.contentType != null) { + if (perRecordContentTypeHolderThreadLocal.get().contentType != null) { record.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)); record.headers().add(header); } catch (Exception e) { @@ -133,7 +134,7 @@ public class KafkaStreamsMessageConversionDelegate { LOG.debug("Could not add content type header"); } } - perRecordContentTypeHolder.unsetContentType(); + perRecordContentTypeHolderThreadLocal.get().unsetContentType(); } }