Class KafkaStreamsFunctionProcessor
java.lang.Object
org.springframework.cloud.stream.binder.kafka.streams.AbstractKafkaStreamsBinderProcessor
org.springframework.cloud.stream.binder.kafka.streams.KafkaStreamsFunctionProcessor
- All Implemented Interfaces:
org.springframework.beans.factory.Aware,org.springframework.beans.factory.BeanFactoryAware,org.springframework.context.ApplicationContextAware
public class KafkaStreamsFunctionProcessor
extends AbstractKafkaStreamsBinderProcessor
implements org.springframework.beans.factory.BeanFactoryAware
- Since:
- 2.2.0
- Author:
- Soby Chacko
-
Field Summary
Fields inherited from class org.springframework.cloud.stream.binder.kafka.streams.AbstractKafkaStreamsBinderProcessor
applicationContext -
Constructor Summary
ConstructorsConstructorDescriptionKafkaStreamsFunctionProcessor(org.springframework.cloud.stream.config.BindingServiceProperties bindingServiceProperties, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, KeyValueSerdeResolver keyValueSerdeResolver, KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate, org.springframework.kafka.core.CleanupConfig cleanupConfig, org.springframework.cloud.stream.function.StreamFunctionProperties streamFunctionProperties, KafkaStreamsBinderConfigurationProperties kafkaStreamsBinderConfigurationProperties, org.springframework.kafka.config.StreamsBuilderFactoryBeanConfigurer customizer, org.springframework.core.env.ConfigurableEnvironment environment) -
Method Summary
Modifier and TypeMethodDescriptionvoidsetBeanFactory(org.springframework.beans.factory.BeanFactory beanFactory) voidsetupFunctionInvokerForKafkaStreams(org.springframework.core.ResolvableType resolvableType, String functionName, KafkaStreamsBindableProxyFactory kafkaStreamsBindableProxyFactory, Method method, org.springframework.core.ResolvableType outputResolvableType, String... composedFunctionNames) This method must be kept stateless.Methods inherited from class org.springframework.cloud.stream.binder.kafka.streams.AbstractKafkaStreamsBinderProcessor
buildStreamsBuilderAndRetrieveConfig, getAutoOffsetReset, getKStream, getValueSerde, handleKTableGlobalKTableInputs, setApplicationContext
-
Constructor Details
-
KafkaStreamsFunctionProcessor
public KafkaStreamsFunctionProcessor(org.springframework.cloud.stream.config.BindingServiceProperties bindingServiceProperties, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, KeyValueSerdeResolver keyValueSerdeResolver, KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate, org.springframework.kafka.core.CleanupConfig cleanupConfig, org.springframework.cloud.stream.function.StreamFunctionProperties streamFunctionProperties, KafkaStreamsBinderConfigurationProperties kafkaStreamsBinderConfigurationProperties, org.springframework.kafka.config.StreamsBuilderFactoryBeanConfigurer customizer, org.springframework.core.env.ConfigurableEnvironment environment)
-
-
Method Details
-
setupFunctionInvokerForKafkaStreams
public void setupFunctionInvokerForKafkaStreams(org.springframework.core.ResolvableType resolvableType, String functionName, KafkaStreamsBindableProxyFactory kafkaStreamsBindableProxyFactory, Method method, org.springframework.core.ResolvableType outputResolvableType, String... composedFunctionNames) This method must be kept stateless. In the case of multiple function beans in an application, isolatedKafkaStreamsBindableProxyFactoryinstances are passed in separately for those functions. If the state is shared between invocations, that will create potential race conditions. Hence, invocations of this method should not be dependent on state modified by a previous invocation.- Parameters:
resolvableType- type of the bindingfunctionName- bean name of the functionkafkaStreamsBindableProxyFactory- bindable proxy factory for the Kafka Streams type
-
setBeanFactory
public void setBeanFactory(org.springframework.beans.factory.BeanFactory beanFactory) throws org.springframework.beans.BeansException - Specified by:
setBeanFactoryin interfaceorg.springframework.beans.factory.BeanFactoryAware- Throws:
org.springframework.beans.BeansException
-