Class KafkaRuleStreamConsumer

java.lang.Object
com.helix.core.stream.kafka.KafkaRuleStreamConsumer
All Implemented Interfaces:
AutoCloseable

public class KafkaRuleStreamConsumer extends Object implements AutoCloseable
Enterprise clustered Kafka stream consumer supporting partition-level parallelism, virtual thread fan-out, dynamic rebalance listeners, and backpressure control via pause/resume.
  • Constructor Details

  • Method Details

    • start

      public void start()
    • subscribe

      public void subscribe(Collection<String> topics)
    • close

      public void close()
      Specified by:
      close in interface AutoCloseable
    • getPausedPartitions

      public Set<org.apache.kafka.common.TopicPartition> getPausedPartitions()
    • getActivePartitions

      public Set<org.apache.kafka.common.TopicPartition> getActivePartitions()
    • getTotalPolled

      public long getTotalPolled()
    • getTotalCompleted

      public long getTotalCompleted()
    • getConsumer

      public org.apache.kafka.clients.consumer.Consumer<String,byte[]> getConsumer()