diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java index b0158c921..50613fa79 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/AbstractKafkaStreamsBinderProcessor.java @@ -492,7 +492,7 @@ public abstract class AbstractKafkaStreamsBinderProcessor implements Application stream = stream.mapValues((value) -> { Object returnValue; String contentType = bindingProperties.getContentType(); - if (value != null && !StringUtils.hasText(contentType)) { + if (value != null && StringUtils.hasText(contentType)) { final Headers headers = headersAtomicReference.get(); final Map headersMap = new HashMap<>(); headers.forEach(header -> headersMap.put(header.key(), header.value()));