Class AbstractKafkaStreamsBinderProcessor
java.lang.Object
org.springframework.cloud.stream.binder.kafka.streams.AbstractKafkaStreamsBinderProcessor
- All Implemented Interfaces:
org.springframework.beans.factory.Aware,org.springframework.context.ApplicationContextAware
- Direct Known Subclasses:
KafkaStreamsFunctionProcessor
public abstract class AbstractKafkaStreamsBinderProcessor
extends Object
implements org.springframework.context.ApplicationContextAware
- Since:
- 3.0.0
- Author:
- Soby Chacko
-
Field Summary
FieldsModifier and TypeFieldDescriptionprotected org.springframework.context.ConfigurableApplicationContext -
Constructor Summary
ConstructorsConstructorDescriptionAbstractKafkaStreamsBinderProcessor(org.springframework.cloud.stream.config.BindingServiceProperties bindingServiceProperties, KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, KeyValueSerdeResolver keyValueSerdeResolver, org.springframework.kafka.core.CleanupConfig cleanupConfig) -
Method Summary
Modifier and TypeMethodDescriptionprotected org.springframework.kafka.config.StreamsBuilderFactoryBeanbuildStreamsBuilderAndRetrieveConfig(String beanNamePostPrefix, org.springframework.context.ApplicationContext applicationContext, String inboundName, KafkaStreamsBinderConfigurationProperties kafkaStreamsBinderConfigurationProperties, org.springframework.kafka.config.StreamsBuilderFactoryBeanConfigurer customizer, org.springframework.core.env.ConfigurableEnvironment environment, org.springframework.cloud.stream.config.BindingProperties bindingProperties) protected org.apache.kafka.streams.Topology.AutoOffsetResetgetAutoOffsetReset(String inboundName, KafkaStreamsConsumerProperties extendedConsumerProperties) protected org.apache.kafka.streams.kstream.KStream<?,?> getKStream(String inboundName, org.springframework.cloud.stream.config.BindingProperties bindingProperties, KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, org.apache.kafka.streams.StreamsBuilder streamsBuilder, org.apache.kafka.common.serialization.Serde<?> keySerde, org.apache.kafka.common.serialization.Serde<?> valueSerde, org.apache.kafka.streams.Topology.AutoOffsetReset autoOffsetReset, boolean firstBuild) protected org.apache.kafka.common.serialization.Serde<?>getValueSerde(String inboundName, KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, org.springframework.core.ResolvableType resolvableType) protected voidhandleKTableGlobalKTableInputs(Object[] arguments, int index, String input, Class<?> parameterType, Object targetBean, org.springframework.kafka.config.StreamsBuilderFactoryBean streamsBuilderFactoryBean, org.apache.kafka.streams.StreamsBuilder streamsBuilder, KafkaStreamsConsumerProperties extendedConsumerProperties, org.apache.kafka.common.serialization.Serde<?> keySerde, org.apache.kafka.common.serialization.Serde<?> valueSerde, org.apache.kafka.streams.Topology.AutoOffsetReset autoOffsetReset, boolean firstBuild) final voidsetApplicationContext(org.springframework.context.ApplicationContext applicationContext)
-
Field Details
-
applicationContext
protected org.springframework.context.ConfigurableApplicationContext applicationContext
-
-
Constructor Details
-
AbstractKafkaStreamsBinderProcessor
public AbstractKafkaStreamsBinderProcessor(org.springframework.cloud.stream.config.BindingServiceProperties bindingServiceProperties, KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, KeyValueSerdeResolver keyValueSerdeResolver, org.springframework.kafka.core.CleanupConfig cleanupConfig)
-
-
Method Details
-
setApplicationContext
public final void setApplicationContext(org.springframework.context.ApplicationContext applicationContext) throws org.springframework.beans.BeansException - Specified by:
setApplicationContextin interfaceorg.springframework.context.ApplicationContextAware- Throws:
org.springframework.beans.BeansException
-
getAutoOffsetReset
protected org.apache.kafka.streams.Topology.AutoOffsetReset getAutoOffsetReset(String inboundName, KafkaStreamsConsumerProperties extendedConsumerProperties) -
handleKTableGlobalKTableInputs
protected void handleKTableGlobalKTableInputs(Object[] arguments, int index, String input, Class<?> parameterType, Object targetBean, org.springframework.kafka.config.StreamsBuilderFactoryBean streamsBuilderFactoryBean, org.apache.kafka.streams.StreamsBuilder streamsBuilder, KafkaStreamsConsumerProperties extendedConsumerProperties, org.apache.kafka.common.serialization.Serde<?> keySerde, org.apache.kafka.common.serialization.Serde<?> valueSerde, org.apache.kafka.streams.Topology.AutoOffsetReset autoOffsetReset, boolean firstBuild) -
buildStreamsBuilderAndRetrieveConfig
protected org.springframework.kafka.config.StreamsBuilderFactoryBean buildStreamsBuilderAndRetrieveConfig(String beanNamePostPrefix, org.springframework.context.ApplicationContext applicationContext, String inboundName, KafkaStreamsBinderConfigurationProperties kafkaStreamsBinderConfigurationProperties, org.springframework.kafka.config.StreamsBuilderFactoryBeanConfigurer customizer, org.springframework.core.env.ConfigurableEnvironment environment, org.springframework.cloud.stream.config.BindingProperties bindingProperties) -
getValueSerde
protected org.apache.kafka.common.serialization.Serde<?> getValueSerde(String inboundName, KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, org.springframework.core.ResolvableType resolvableType) -
getKStream
protected org.apache.kafka.streams.kstream.KStream<?,?> getKStream(String inboundName, org.springframework.cloud.stream.config.BindingProperties bindingProperties, KafkaStreamsConsumerProperties kafkaStreamsConsumerProperties, org.apache.kafka.streams.StreamsBuilder streamsBuilder, org.apache.kafka.common.serialization.Serde<?> keySerde, org.apache.kafka.common.serialization.Serde<?> valueSerde, org.apache.kafka.streams.Topology.AutoOffsetReset autoOffsetReset, boolean firstBuild)
-