From e7c1c7f7690a790aac4f83469b61bfad60cf9b15 Mon Sep 17 00:00:00 2001 From: LazroLeader <135716376+LazroLeader@users.noreply.github.com> Date: Fri, 16 Feb 2024 00:58:07 +0400 Subject: [PATCH] Fixing concurrency issue in Kafka Streams binder * Thread Safety Issue in serializeOnOutbound Method of KafkaStreamsMessageConversionDelegate * Wrapped perRecordContentTypeHolder with ThreadLocal * update year and author --- .../KafkaStreamsMessageConversionDelegate.java | 13 +++++++------ 1 file changed, 7 insertions(+), 6 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 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(); } }