From 98431ed8a023ced69ca6821c0fd9e3299e920444 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 9 Oct 2019 00:28:54 -0400 Subject: [PATCH] Fix spurious warnings from InteractiveQueryService Set applicationId properly in functions with multiple inputs --- .../streams/InteractiveQueryService.java | 23 ++++++++++++------- .../KafkaStreamsFunctionProcessor.java | 3 +++ ...StreamListenerSetupMethodOrchestrator.java | 2 ++ 3 files changed, 20 insertions(+), 8 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/InteractiveQueryService.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/InteractiveQueryService.java index 3dc7e201a..9256d2893 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/InteractiveQueryService.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/InteractiveQueryService.java @@ -16,8 +16,10 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import java.util.Iterator; import java.util.Map; import java.util.Optional; +import java.util.Set; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -84,19 +86,24 @@ public class InteractiveQueryService { retryTemplate.setRetryPolicy(retryPolicy); return retryTemplate.execute(context -> { - T store; - for (KafkaStreams kafkaStream : InteractiveQueryService.this.kafkaStreamsRegistry.getKafkaStreams()) { + T store = null; + + final Set kafkaStreams = InteractiveQueryService.this.kafkaStreamsRegistry.getKafkaStreams(); + final Iterator iterator = kafkaStreams.iterator(); + Throwable throwable = null; + while (iterator.hasNext()) { try { - store = kafkaStream.store(storeName, storeType); - if (store != null) { - return store; - } + store = iterator.next().store(storeName, storeType); } catch (InvalidStateStoreException e) { - LOG.warn("Error when retrieving state store: " + storeName, e); + // pass through.. + throwable = e; } } - throw new IllegalStateException("Error when retrieving state store: " + storeName); + if (store != null) { + return store; + } + throw new IllegalStateException("Error when retrieving state store: j " + storeName, throwable); }); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java index 4875cc828..6408c564d 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsFunctionProcessor.java @@ -35,6 +35,7 @@ import org.apache.commons.logging.LogFactory; import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.StreamsBuilder; +import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.Topology; import org.apache.kafka.streams.kstream.KStream; @@ -296,8 +297,10 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro StreamsBuilderFactoryBean streamsBuilderFactoryBean = this.methodStreamsBuilderFactoryBeanMap.get(functionName); StreamsBuilder streamsBuilder = streamsBuilderFactoryBean.getObject(); + final String applicationId = streamsBuilderFactoryBean.getStreamsConfiguration().getProperty(StreamsConfig.APPLICATION_ID_CONFIG); KafkaStreamsConsumerProperties extendedConsumerProperties = this.kafkaStreamsExtendedBindingProperties.getExtendedConsumerProperties(input); + extendedConsumerProperties.setApplicationId(applicationId); //get state store spec Serde keySerde = this.keyValueSerdeResolver.getInboundKeySerde(extendedConsumerProperties, stringResolvableTypeMap.get(input)); diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java index 094f6f1a1..e07f5da36 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsStreamListenerSetupMethodOrchestrator.java @@ -251,8 +251,10 @@ class KafkaStreamsStreamListenerSetupMethodOrchestrator extends AbstractKafkaStr StreamsBuilderFactoryBean streamsBuilderFactoryBean = this.methodStreamsBuilderFactoryBeanMap .get(method); StreamsBuilder streamsBuilder = streamsBuilderFactoryBean.getObject(); + final String applicationId = streamsBuilderFactoryBean.getStreamsConfiguration().getProperty(StreamsConfig.APPLICATION_ID_CONFIG); KafkaStreamsConsumerProperties extendedConsumerProperties = this.kafkaStreamsExtendedBindingProperties .getExtendedConsumerProperties(inboundName); + extendedConsumerProperties.setApplicationId(applicationId); // get state store spec KafkaStreamsStateStoreProperties spec = buildStateStoreSpec(method);