Package com.helix.core.stream.kafka
Class KafkaRuleStreamConsumer
java.lang.Object
com.helix.core.stream.kafka.KafkaRuleStreamConsumer
- All Implemented Interfaces:
AutoCloseable
Enterprise clustered Kafka stream consumer supporting partition-level parallelism,
virtual thread fan-out, dynamic rebalance listeners, and backpressure control via pause/resume.
-
Nested Class Summary
Nested Classes -
Constructor Summary
ConstructorsConstructorDescriptionKafkaRuleStreamConsumer(KafkaStreamConfig config, KafkaRuleStreamConsumer.RecordEvaluator evaluator) KafkaRuleStreamConsumer(KafkaStreamConfig config, org.apache.kafka.clients.consumer.Consumer<String, byte[]> consumer, KafkaRuleStreamConsumer.RecordEvaluator evaluator) -
Method Summary
-
Constructor Details
-
KafkaRuleStreamConsumer
public KafkaRuleStreamConsumer(KafkaStreamConfig config, KafkaRuleStreamConsumer.RecordEvaluator evaluator) -
KafkaRuleStreamConsumer
public KafkaRuleStreamConsumer(KafkaStreamConfig config, org.apache.kafka.clients.consumer.Consumer<String, byte[]> consumer, KafkaRuleStreamConsumer.RecordEvaluator evaluator)
-
-
Method Details
-
start
public void start() -
subscribe
-
close
public void close()- Specified by:
closein interfaceAutoCloseable
-
getPausedPartitions
-
getActivePartitions
-
getTotalPolled
public long getTotalPolled() -
getTotalCompleted
public long getTotalCompleted() -
getConsumer
-