From 498998f1a4a1f6d4091e01c1e7c802d60aeadba4 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 1 May 2019 21:12:40 +0200 Subject: [PATCH] polishing Resolves #643 --- .../KafkaStreamsFunctionProcessor.java | 21 +++++++------------ 1 file changed, 8 insertions(+), 13 deletions(-) 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 ec7f15ad5..baa7697d8 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 @@ -148,23 +148,19 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { return resolvableTypeMap; } - @SuppressWarnings("unchecked") + @SuppressWarnings({ "unchecked", "rawtypes" }) public void orchestrateFunctionInvoking(ResolvableType resolvableType, String functionName) { final Map stringResolvableTypeMap = buildTypeMap(resolvableType); Object[] adaptedInboundArguments = adaptAndRetrieveInboundArguments(stringResolvableTypeMap, functionName); try { if (resolvableType.getRawClass() != null && resolvableType.getRawClass().equals(Consumer.class)) { - //TOOD: Investigate why looking up by Consumer returns null - Consumer consumer = functionCatalog.lookup(Consumer.class, functionName); - if (consumer == null) { - FluxedConsumer fluxedConsumer = functionCatalog.lookup(FluxedConsumer.class, functionName); - Assert.isTrue(fluxedConsumer != null, - "No corresponding consumer beans found in the catalog"); - Object target = fluxedConsumer.getTarget(); - if (Consumer.class.isAssignableFrom(target.getClass())) { - consumer = (Consumer) target; - } - } + FluxedConsumer fluxedConsumer = functionCatalog.lookup(FluxedConsumer.class, functionName); + Assert.isTrue(fluxedConsumer != null, + "No corresponding consumer beans found in the catalog"); + Object target = fluxedConsumer.getTarget(); + + Consumer consumer = Consumer.class.isAssignableFrom(target.getClass()) ? (Consumer) target : null; + if (consumer != null) { consumer.accept(adaptedInboundArguments[0]); } @@ -257,7 +253,6 @@ public class KafkaStreamsFunctionProcessor implements ApplicationContextAware { KafkaStreamsConsumerProperties extendedConsumerProperties = this.kafkaStreamsExtendedBindingProperties.getExtendedConsumerProperties(input); //get state store spec - //KafkaStreamsStateStoreProperties spec = buildStateStoreSpec(method); Serde keySerde = this.keyValueSerdeResolver.getInboundKeySerde(extendedConsumerProperties); Serde valueSerde = this.keyValueSerdeResolver.getInboundValueSerde( bindingProperties.getConsumer(), extendedConsumerProperties);