Class KafkaStreamConfig

java.lang.Object
com.helix.core.stream.kafka.KafkaStreamConfig

public class KafkaStreamConfig extends Object
Configuration for clustered Kafka rule streaming ingestion and result publishing.
  • Field Details

    • DEFAULT_BOOTSTRAP_SERVERS

      public static final String DEFAULT_BOOTSTRAP_SERVERS
      See Also:
    • DEFAULT_GROUP_ID

      public static final String DEFAULT_GROUP_ID
      See Also:
    • DEFAULT_AUTO_OFFSET_RESET

      public static final String DEFAULT_AUTO_OFFSET_RESET
      See Also:
    • DEFAULT_POLL_TIMEOUT

      public static final Duration DEFAULT_POLL_TIMEOUT
    • DEFAULT_MAX_POLL_RECORDS

      public static final int DEFAULT_MAX_POLL_RECORDS
      See Also:
    • DEFAULT_MAX_IN_FLIGHT_PER_PARTITION

      public static final int DEFAULT_MAX_IN_FLIGHT_PER_PARTITION
      See Also:
    • DEFAULT_LOW_WATERMARK_PER_PARTITION

      public static final int DEFAULT_LOW_WATERMARK_PER_PARTITION
      See Also:
  • Constructor Details

  • Method Details

    • builder

      public static KafkaStreamConfig.Builder builder()
    • defaultConfig

      public static KafkaStreamConfig defaultConfig()
    • toConsumerProperties

      public Properties toConsumerProperties()
    • toProducerProperties

      public Properties toProducerProperties()
    • getBootstrapServers

      public String getBootstrapServers()
    • getGroupId

      public String getGroupId()
    • getInputTopics

      public Set<String> getInputTopics()
    • getOutputTopic

      public String getOutputTopic()
    • getAutoOffsetReset

      public String getAutoOffsetReset()
    • isEnableAutoCommit

      public boolean isEnableAutoCommit()
    • getCommitMode

      public CommitMode getCommitMode()
    • getPollTimeout

      public Duration getPollTimeout()
    • getMaxPollRecords

      public int getMaxPollRecords()
    • getMaxInFlightPerPartition

      public int getMaxInFlightPerPartition()
    • getLowWatermarkPerPartition

      public int getLowWatermarkPerPartition()
    • getCustomConsumerProps

      public Map<String,Object> getCustomConsumerProps()
    • getCustomProducerProps

      public Map<String,Object> getCustomProducerProps()