diff --git a/pom.xml b/pom.xml index 5833cbee4..9fd6d8b0d 100644 --- a/pom.xml +++ b/pom.xml @@ -124,18 +124,6 @@ maven-antrun-plugin 1.7 - - org.apache.maven.plugins - maven-checkstyle-plugin - 2.17 - - - com.puppycrawl.tools - checkstyle - 7.1 - - - org.apache.maven.plugins maven-javadoc-plugin @@ -159,26 +147,15 @@ org.springframework.cloud - spring-cloud-build-tools - ${project.parent.version} + spring-cloud-stream-tools + ${spring-cloud-stream.version} - - - checkstyle-validation - validate - - checkstyle.xml - UTF-8 - true - true - true - - - check - - - + + checkstyle.xml + checkstyle-header.txt + true + diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/JaasLoginModuleConfiguration.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/JaasLoginModuleConfiguration.java index 33e9865f1..f82ec51f3 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/JaasLoginModuleConfiguration.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/JaasLoginModuleConfiguration.java @@ -18,6 +18,7 @@ package org.springframework.cloud.stream.binder.kafka.properties; import java.util.HashMap; import java.util.Map; + import javax.security.auth.login.AppConfigurationEntry; import org.springframework.util.Assert; diff --git a/spring-cloud-stream-binder-kafka-test/pom.xml b/spring-cloud-stream-binder-kafka-test/pom.xml index 798f21be3..b80169b03 100644 --- a/spring-cloud-stream-binder-kafka-test/pom.xml +++ b/spring-cloud-stream-binder-kafka-test/pom.xml @@ -4,10 +4,7 @@ org.springframework.cloud spring-cloud-stream-binder-kafka-parent - 2.0.0.M1 -======= 2.0.0.BUILD-SNAPSHOT ->>>>>>> Set version to 2.0.0.BUILD-SNAPSHOT spring-cloud-stream-binder-kafka-test Spring Cloud Stream Kafka Binder Tests @@ -66,13 +63,6 @@ org.springframework.integration spring-integration-jmx - - org.springframework.cloud - spring-cloud-stream-binder-kafka-0.10.2-test - ${project.version} - test-jar - test - org.springframework.cloud spring-cloud-stream-binder-test diff --git a/spring-cloud-stream-binder-kafka/pom.xml b/spring-cloud-stream-binder-kafka/pom.xml index c8a939fef..d0d143a16 100644 --- a/spring-cloud-stream-binder-kafka/pom.xml +++ b/spring-cloud-stream-binder-kafka/pom.xml @@ -28,10 +28,6 @@ org.springframework.cloud spring-cloud-stream - - org.springframework.cloud - spring-cloud-stream-codec - org.springframework.boot spring-boot-autoconfigure @@ -81,18 +77,6 @@ ${spring-cloud-stream.version} test - - io.confluent - kafka-avro-serializer - 3.2.1 - test - - - io.confluent - kafka-schema-registry - 3.2.1 - test - diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java index 9240362d6..b6abaeb28 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetrics.java @@ -1,11 +1,11 @@ /* - * Copyright 2017 the original author or authors. + * Copyright 2016-2017 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. * You may obtain a copy of the License at * - * http://www.apache.org/licenses/LICENSE-2.0 + * http://www.apache.org/licenses/LICENSE-2.0 * * Unless required by applicable law or agreed to in writing, software * distributed under the License is distributed on an "AS IS" BASIS, @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.cloud.stream.binder.kafka; import java.util.HashMap; @@ -22,14 +23,14 @@ import java.util.Map; import io.micrometer.core.instrument.MeterRegistry; import io.micrometer.core.instrument.binder.MeterBinder; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.serialization.ByteArrayDeserializer; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.kafka.core.ConsumerFactory; @@ -43,7 +44,7 @@ import org.springframework.util.ObjectUtils; */ public class KafkaBinderMetrics implements MeterBinder { - private final static Logger LOG = LoggerFactory.getLogger(KafkaBinderMetrics.class); + private final static Log LOG = LogFactory.getLog(KafkaBinderMetrics.class); static final String METRIC_PREFIX = "spring.cloud.stream.binder.kafka"; @@ -120,4 +121,4 @@ public class KafkaBinderMetrics implements MeterBinder { return new DefaultKafkaConsumerFactory<>(props); } -} \ No newline at end of file +} diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index e3a7dcbeb..ff21fb643 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -39,7 +39,6 @@ import org.apache.kafka.common.utils.Utils; import org.springframework.beans.factory.DisposableBean; import org.springframework.cloud.stream.binder.AbstractMessageChannelBinder; -import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.binder.BinderHeaders; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedProducerProperties; @@ -85,7 +84,7 @@ import org.springframework.util.concurrent.ListenableFuture; import org.springframework.util.concurrent.ListenableFutureCallback; /** - * A {@link Binder} that uses Kafka as the underlying middleware. + * A {@link org.springframework.cloud.stream.binder.Binder} that uses Kafka as the underlying middleware. * * @author Eric Bottard * @author Marius Bogoevici diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AbstractKafkaTestBinder.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AbstractKafkaTestBinder.java index b8bd62e31..dd25bd4c6 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AbstractKafkaTestBinder.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AbstractKafkaTestBinder.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.cloud.stream.binder.kafka; import org.springframework.cloud.stream.binder.AbstractTestBinder; diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderAutoConfigurationPropertiesTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderAutoConfigurationPropertiesTest.java index ca9e6333f..2b4e89460 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderAutoConfigurationPropertiesTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderAutoConfigurationPropertiesTest.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.cloud.stream.binder.kafka; import java.lang.reflect.Field; diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java index 7a85a8b1b..ee40aaac7 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationPropertiesTest.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.cloud.stream.binder.kafka; import java.lang.reflect.Field; diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationTest.java index 6ff911e9b..23ba781a9 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationTest.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.cloud.stream.binder.kafka; import java.lang.reflect.Field; diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicatorTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicatorTest.java index 72b948af3..985543068 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicatorTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicatorTest.java @@ -13,11 +13,8 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.springframework.cloud.stream.binder.kafka; -import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.BDDMockito.given; -import static org.mockito.Mockito.verify; +package org.springframework.cloud.stream.binder.kafka; import java.util.ArrayList; import java.util.HashMap; @@ -38,6 +35,8 @@ import org.springframework.boot.actuate.health.Health; import org.springframework.boot.actuate.health.Status; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; +import static org.assertj.core.api.Assertions.assertThat; + /** * @author Barry Commins * @author Gary Russell @@ -62,8 +61,8 @@ public class KafkaBinderHealthIndicatorTest { @Before public void setup() { MockitoAnnotations.initMocks(this); - given(consumerFactory.createConsumer()).willReturn(consumer); - given(binder.getTopicsInUse()).willReturn(topicsInUse); + org.mockito.BDDMockito.given(consumerFactory.createConsumer()).willReturn(consumer); + org.mockito.BDDMockito.given(binder.getTopicsInUse()).willReturn(topicsInUse); this.indicator = new KafkaBinderHealthIndicator(binder, consumerFactory); this.indicator.setTimeout(10); } @@ -72,17 +71,17 @@ public class KafkaBinderHealthIndicatorTest { public void kafkaBinderIsUp() { final List partitions = partitions(new Node(0, null, 0)); topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group", partitions)); - given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); + org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); Health health = indicator.health(); assertThat(health.getStatus()).isEqualTo(Status.UP); - verify(this.consumer).close(); + org.mockito.Mockito.verify(this.consumer).close(); } @Test public void kafkaBinderIsDown() { final List partitions = partitions(new Node(-1, null, 0)); topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group", partitions)); - given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); + org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); Health health = indicator.health(); assertThat(health.getStatus()).isEqualTo(Status.DOWN); } @@ -91,7 +90,7 @@ public class KafkaBinderHealthIndicatorTest { public void kafkaBinderDoesNotAnswer() { final List partitions = partitions(new Node(-1, null, 0)); topicsInUse.put(TEST_TOPIC, new KafkaMessageChannelBinder.TopicInformation("group", partitions)); - given(consumer.partitionsFor(TEST_TOPIC)).willAnswer(new Answer() { + org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willAnswer(new Answer() { @Override public Object answer(InvocationOnMock invocation) throws Throwable { diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java index 58b9cbd89..e546cef2c 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderMetricsTest.java @@ -1,5 +1,5 @@ /* - * Copyright 2017 the original author or authors. + * Copyright 2016-2017 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. @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.cloud.stream.binder.kafka; import java.util.ArrayList; @@ -36,12 +37,7 @@ import org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder.T import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; -import static java.util.Collections.singletonMap; import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.BDDMockito.given; -import static org.mockito.Matchers.any; -import static org.mockito.Matchers.anyCollectionOf; -import static org.springframework.cloud.stream.binder.kafka.KafkaBinderMetrics.METRIC_PREFIX; /** * @author Henryk Konsek @@ -71,22 +67,22 @@ public class KafkaBinderMetricsTest { @Before public void setup() { MockitoAnnotations.initMocks(this); - given(consumerFactory.createConsumer()).willReturn(consumer); - given(binder.getTopicsInUse()).willReturn(topicsInUse); + org.mockito.BDDMockito.given(consumerFactory.createConsumer()).willReturn(consumer); + org.mockito.BDDMockito.given(binder.getTopicsInUse()).willReturn(topicsInUse); metrics = new KafkaBinderMetrics(binder, kafkaBinderConfigurationProperties, consumerFactory); - given(consumer.endOffsets(anyCollectionOf(TopicPartition.class))) - .willReturn(singletonMap(new TopicPartition(TEST_TOPIC, 0), 1000L)); + org.mockito.BDDMockito.given(consumer.endOffsets(org.mockito.Matchers.anyCollectionOf(TopicPartition.class))) + .willReturn(java.util.Collections.singletonMap(new TopicPartition(TEST_TOPIC, 0), 1000L)); } @Test public void shouldIndicateLag() { - given(consumer.committed(any(TopicPartition.class))).willReturn(new OffsetAndMetadata(500)); + org.mockito.BDDMockito.given(consumer.committed(org.mockito.Matchers.any(TopicPartition.class))).willReturn(new OffsetAndMetadata(500)); List partitions = partitions(new Node(0, null, 0)); topicsInUse.put(TEST_TOPIC, new TopicInformation("group", partitions)); - given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); + org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); - MeterRegistry.Search group = meterRegistry.find(String.format("%s.%s.%s.lag", METRIC_PREFIX, "group", TEST_TOPIC)); + MeterRegistry.Search group = meterRegistry.find(String.format("%s.%s.%s.lag", KafkaBinderMetrics.METRIC_PREFIX, "group", TEST_TOPIC)); assertThat(group.gauge().get().value()).isEqualTo(500.0); } @@ -95,14 +91,14 @@ public class KafkaBinderMetricsTest { Map endOffsets = new HashMap<>(); endOffsets.put(new TopicPartition(TEST_TOPIC, 0), 1000L); endOffsets.put(new TopicPartition(TEST_TOPIC, 1), 1000L); - given(consumer.endOffsets(anyCollectionOf(TopicPartition.class))).willReturn(endOffsets); - given(consumer.committed(any(TopicPartition.class))).willReturn(new OffsetAndMetadata(500)); + org.mockito.BDDMockito.given(consumer.endOffsets(org.mockito.Matchers.anyCollectionOf(TopicPartition.class))).willReturn(endOffsets); + org.mockito.BDDMockito.given(consumer.committed(org.mockito.Matchers.any(TopicPartition.class))).willReturn(new OffsetAndMetadata(500)); List partitions = partitions(new Node(0, null, 0), new Node(0, null, 0)); topicsInUse.put(TEST_TOPIC, new TopicInformation("group", partitions)); - given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); + org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); - MeterRegistry.Search group = meterRegistry.find(String.format("%s.%s.%s.lag", METRIC_PREFIX, "group", TEST_TOPIC)); + MeterRegistry.Search group = meterRegistry.find(String.format("%s.%s.%s.lag", KafkaBinderMetrics.METRIC_PREFIX, "group", TEST_TOPIC)); assertThat(group.gauge().get().value()).isEqualTo(1000.0); } @@ -110,10 +106,10 @@ public class KafkaBinderMetricsTest { public void shouldIndicateFullLagForNotCommittedGroups() { List partitions = partitions(new Node(0, null, 0)); topicsInUse.put(TEST_TOPIC, new TopicInformation("group", partitions)); - given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); + org.mockito.BDDMockito.given(consumer.partitionsFor(TEST_TOPIC)).willReturn(partitions); metrics.bindTo(meterRegistry); assertThat(meterRegistry.getMeters()).hasSize(1); - MeterRegistry.Search group = meterRegistry.find(String.format("%s.%s.%s.lag", METRIC_PREFIX, "group", TEST_TOPIC)); + MeterRegistry.Search group = meterRegistry.find(String.format("%s.%s.%s.lag", KafkaBinderMetrics.METRIC_PREFIX, "group", TEST_TOPIC)); assertThat(group.gauge().get().value()).isEqualTo(1000.0); } diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/bootstrap/KafkaBinderBootstrapTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/bootstrap/KafkaBinderBootstrapTest.java index ad672c892..f9d7e13c2 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/bootstrap/KafkaBinderBootstrapTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/bootstrap/KafkaBinderBootstrapTest.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.cloud.stream.binder.kafka.bootstrap; import org.junit.ClassRule; diff --git a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamBoundElementFactory.java b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamBoundElementFactory.java index 7a6cc376f..701326a4d 100644 --- a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamBoundElementFactory.java +++ b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamBoundElementFactory.java @@ -60,16 +60,18 @@ public class KStreamBoundElementFactory extends AbstractBindingTargetFactory stream = kStreamBuilder.stream(bindingServiceProperties.getBindingDestination(name)); stream = stream.map((key, value) -> { - + KeyValue keyValue; BindingProperties bindingProperties = bindingServiceProperties.getBindingProperties(name); String contentType = bindingProperties.getContentType(); if (!StringUtils.isEmpty(contentType)) { - Message message = MessageBuilder.withPayload(value) .setHeader(MessageHeaders.CONTENT_TYPE, contentType).build(); - return new KeyValue<>(key, message); + keyValue = new KeyValue<>(key, message); } - return new KeyValue<>(key, value); + else { + keyValue = new KeyValue<>(key, value); + } + return keyValue; }); return stream; } @@ -87,12 +89,12 @@ public class KStreamBoundElementFactory extends AbstractBindingTargetFactory delegate); } - static class KStreamWrapperHandler implements KStreamWrapper, MethodInterceptor { + private static class KStreamWrapperHandler implements KStreamWrapper, MethodInterceptor { private KStream delegate; @@ -100,7 +102,7 @@ public class KStreamBoundElementFactory extends AbstractBindingTargetFactory { + KeyValue keyValue; if (valueClass.isAssignableFrom(o2.getClass())) { - return new KeyValue<>(o, o2); + keyValue = new KeyValue<>(o, o2); } else if (o2 instanceof Message) { if (valueClass.isAssignableFrom(((Message) o2).getPayload().getClass())) { - return new KeyValue<>(o, ((Message) o2).getPayload()); + keyValue = new KeyValue<>(o, ((Message) o2).getPayload()); + } + else { + keyValue = new KeyValue<>(o, messageConverter.fromMessage((Message) o2, valueClass)); } - return new KeyValue<>(o, messageConverter.fromMessage((Message) o2, valueClass)); } else if(o2 instanceof String || o2 instanceof byte[]) { Message message = MessageBuilder.withPayload(o2).build(); - return new KeyValue<>(o, messageConverter.fromMessage(message, valueClass)); + keyValue = new KeyValue<>(o, messageConverter.fromMessage(message, valueClass)); } else { - return new KeyValue<>(o, o2); + keyValue = new KeyValue<>(o, o2); } + return keyValue; }); } diff --git a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamStreamListenerResultAdapter.java b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamStreamListenerResultAdapter.java index e1585ff3c..843fd1b98 100644 --- a/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamStreamListenerResultAdapter.java +++ b/spring-cloud-stream-binder-kstream/src/main/java/org/springframework/cloud/stream/binder/kstream/KStreamStreamListenerResultAdapter.java @@ -40,12 +40,14 @@ public class KStreamStreamListenerResultAdapter implements StreamListenerResultA @SuppressWarnings("unchecked") public Closeable adapt(KStream streamListenerResult, KStreamBoundElementFactory.KStreamWrapper boundElement) { boundElement.wrap(streamListenerResult.map((k, v) -> { + KeyValue keyValue; if (v instanceof Message) { - return new KeyValue<>(k, v); + keyValue = new KeyValue<>(k, v); } else { - return new KeyValue<>(k, MessageBuilder.withPayload(v).build()); + keyValue = new KeyValue<>(k, MessageBuilder.withPayload(v).build()); } + return keyValue; })); return new NoOpCloseable(); } diff --git a/spring-cloud-stream-binder-kstream/src/test/java/org/springframework/cloud/stream/binder/kstream/KStreamInteractiveQueryIntegrationTests.java b/spring-cloud-stream-binder-kstream/src/test/java/org/springframework/cloud/stream/binder/kstream/KStreamInteractiveQueryIntegrationTests.java index 68d249f8a..775987371 100644 --- a/spring-cloud-stream-binder-kstream/src/test/java/org/springframework/cloud/stream/binder/kstream/KStreamInteractiveQueryIntegrationTests.java +++ b/spring-cloud-stream-binder-kstream/src/test/java/org/springframework/cloud/stream/binder/kstream/KStreamInteractiveQueryIntegrationTests.java @@ -13,6 +13,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.cloud.stream.binder.kstream; import java.util.Map;