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

    Fields
    Modifier and Type
    Field
    Description
    protected org.springframework.context.ConfigurableApplicationContext
     
  • Constructor Summary

    Constructors
    Constructor
    Description
    AbstractKafkaStreamsBinderProcessor(org.springframework.cloud.stream.config.BindingServiceProperties bindingServiceProperties, KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, KeyValueSerdeResolver keyValueSerdeResolver, org.springframework.kafka.core.CleanupConfig cleanupConfig)
     
  • Method Summary

    Modifier and Type
    Method
    Description
    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)
     
    protected org.apache.kafka.streams.Topology.AutoOffsetReset
    getAutoOffsetReset(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 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)
     
    final void
    setApplicationContext(org.springframework.context.ApplicationContext applicationContext)
     

    Methods inherited from class java.lang.Object

    clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
  • Field Details

    • applicationContext

      protected org.springframework.context.ConfigurableApplicationContext applicationContext
  • Constructor Details

  • Method Details

    • setApplicationContext

      public final void setApplicationContext(org.springframework.context.ApplicationContext applicationContext) throws org.springframework.beans.BeansException
      Specified by:
      setApplicationContext in interface org.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)