diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java index 406a1f100..1e08f4012 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java @@ -19,8 +19,6 @@ package org.springframework.cloud.stream.binder.kafka.properties; import java.util.HashMap; import java.util.Map; -import org.springframework.boot.context.properties.DeprecatedConfigurationProperty; - /** * Extended consumer properties for Kafka binder. * 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 1d1b4999f..4807f3d78 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 @@ -51,6 +51,7 @@ import org.springframework.cloud.stream.binder.kafka.streams.properties.KafkaStr import org.springframework.cloud.stream.binding.StreamListenerErrorMessages; import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.BindingServiceProperties; +import org.springframework.cloud.stream.function.FunctionConstants; import org.springframework.cloud.stream.function.StreamFunctionProperties; import org.springframework.core.ResolvableType; import org.springframework.kafka.config.StreamsBuilderFactoryBean; @@ -267,7 +268,7 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro } else { for (int i = 0; i < outputs; i++) { - outputBindingNames.add(String.format("%s-%s-%d", functionName, KafkaStreamsBindableProxyFactory.DEFAULT_OUTPUT_SUFFIX, i)); + outputBindingNames.add(String.format("%s-%s-%d", functionName, FunctionConstants.DEFAULT_OUTPUT_SUFFIX, i)); } } return outputBindingNames; @@ -313,7 +314,7 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro //wrap the proxy created during the initial target type binding with real object (KStream) kStreamWrapper.wrap((KStream) stream); - this.kafkaStreamsBindingInformationCatalogue.addKeySerde((KStream) kStreamWrapper, keySerde); + this.kafkaStreamsBindingInformationCatalogue.addKeySerde((KStream) kStreamWrapper, keySerde); this.kafkaStreamsBindingInformationCatalogue.addStreamBuilderFactory(streamsBuilderFactoryBean); if (KStream.class.isAssignableFrom(stringResolvableTypeMap.get(input).getRawClass())) { @@ -321,7 +322,7 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro (stringResolvableTypeMap.get(input).getGeneric(1).getRawClass() != null) ? (stringResolvableTypeMap.get(input).getGeneric(1).getRawClass()) : Object.class; if (this.kafkaStreamsBindingInformationCatalogue.isUseNativeDecoding( - (KStream) kStreamWrapper)) { + (KStream) kStreamWrapper)) { arguments[i] = stream; } else { @@ -352,13 +353,6 @@ public class KafkaStreamsFunctionProcessor extends AbstractKafkaStreamsBinderPro return arguments; } - private KStream getkStream(String inboundName, - BindingProperties bindingProperties, - StreamsBuilder streamsBuilder, - Serde keySerde, Serde valueSerde, Topology.AutoOffsetReset autoOffsetReset) { - return getKStream(inboundName, bindingProperties, streamsBuilder, keySerde, valueSerde, autoOffsetReset); - } - @Override public void setBeanFactory(BeanFactory beanFactory) throws BeansException { this.beanFactory = beanFactory; diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBindableProxyFactory.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBindableProxyFactory.java index 051a440b6..e5c5fa1c6 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBindableProxyFactory.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/function/KafkaStreamsBindableProxyFactory.java @@ -27,8 +27,6 @@ import java.util.function.BiFunction; import java.util.function.Consumer; import java.util.function.Function; -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; import org.apache.kafka.streams.kstream.GlobalKTable; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.KTable; @@ -41,8 +39,8 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.beans.factory.support.RootBeanDefinition; import org.springframework.cloud.stream.binding.AbstractBindableProxyFactory; -import org.springframework.cloud.stream.binding.BindableProxyFactory; import org.springframework.cloud.stream.binding.BoundTargetHolder; +import org.springframework.cloud.stream.function.FunctionConstants; import org.springframework.cloud.stream.function.StreamFunctionProperties; import org.springframework.core.ResolvableType; import org.springframework.util.Assert; @@ -69,15 +67,6 @@ import org.springframework.util.CollectionUtils; */ public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFactory implements InitializingBean, BeanFactoryAware { - /** - * Default output binding name. Output binding may occur later on in the function invoker (outside of this class), - * thus making this field part of the API. - */ - public static final String DEFAULT_OUTPUT_SUFFIX = "out"; - private static final String DEFAULT_INPUT_SUFFIX = "in"; - - private static Log log = LogFactory.getLog(BindableProxyFactory.class); - @Autowired private StreamFunctionProperties streamFunctionProperties; @@ -150,7 +139,7 @@ public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFacto outputBinding = "output"; } else { - outputBinding = String.format("%s-%s-0", this.functionName, DEFAULT_OUTPUT_SUFFIX); + outputBinding = String.format("%s-%s-0", this.functionName, FunctionConstants.DEFAULT_OUTPUT_SUFFIX); } } Assert.isTrue(outputBinding != null, "output binding is not inferred."); @@ -205,14 +194,14 @@ public class KafkaStreamsBindableProxyFactory extends AbstractBindableProxyFacto inputs.add("input"); } else { - inputs.add(String.format("%s-%s-0", this.functionName, DEFAULT_INPUT_SUFFIX)); + inputs.add(String.format("%s-%s-0", this.functionName, FunctionConstants.DEFAULT_INPUT_SUFFIX)); } return inputs; } else { int i = 0; while (i < numberOfInputs) { - inputs.add(String.format("%s-%s-%d", this.functionName, DEFAULT_INPUT_SUFFIX, i++)); + inputs.add(String.format("%s-%s-%d", this.functionName, FunctionConstants.DEFAULT_INPUT_SUFFIX, i++)); } return inputs; }