Class KafkaStreamEngine

java.lang.Object
com.helix.core.stream.kafka.KafkaStreamEngine
All Implemented Interfaces:
StreamBridge, AutoCloseable

public class KafkaStreamEngine extends Object implements StreamBridge
Enterprise clustered Kafka streaming engine implementing StreamBridge. Integrates partitioned Kafka consumption with virtual-thread fan-out, rule execution, dynamic offset commitment, backpressure control, and result publication.
  • Constructor Details

    • KafkaStreamEngine

      public KafkaStreamEngine(KafkaStreamConfig config)
    • KafkaStreamEngine

      public KafkaStreamEngine(KafkaStreamConfig config, org.apache.kafka.clients.consumer.Consumer<String,byte[]> consumer, org.apache.kafka.clients.producer.Producer<String,byte[]> producer)
  • Method Details

    • start

      public void start()
    • subscribe

      public void subscribe(String topic, RuleStreamHandler handler)
      Description copied from interface: StreamBridge
      Subscribes a functional rule stream handler to a given topic.
      Specified by:
      subscribe in interface StreamBridge
      Parameters:
      topic - stream topic
      handler - handler invoked for incoming events on the topic
    • subscribeListener

      public void subscribeListener(String topic, StreamListener listener)
      Description copied from interface: StreamBridge
      Subscribes an asynchronous result listener to a given topic.
      Specified by:
      subscribeListener in interface StreamBridge
      Parameters:
      topic - stream topic
      listener - listener invoked when events on the topic complete evaluation
    • registerRule

      public void registerRule(String topicOrRuleName, CompiledRule rule)
    • publish

      public CompletableFuture<StreamResult> publish(RuleEvent event)
      Description copied from interface: StreamBridge
      Publishes a rule evaluation event into the stream and returns a CompletableFuture for the result.
      Specified by:
      publish in interface StreamBridge
      Parameters:
      event - the rule event to publish
      Returns:
      CompletableFuture completing with the evaluation result
    • publish

      public <T> CompletableFuture<StreamResult> publish(StreamRecord<T> record)
      Description copied from interface: StreamBridge
      Publishes a generic stream record into the stream and returns a CompletableFuture for the result.
      Specified by:
      publish in interface StreamBridge
      Type Parameters:
      T - payload type
      Parameters:
      record - the stream record to publish
      Returns:
      CompletableFuture completing with the evaluation result
    • tryPublish

      public boolean tryPublish(RuleEvent event, StreamListener listener)
      Description copied from interface: StreamBridge
      Fast-path non-blocking publish of a rule event with a completion listener, avoiding Future allocation.
      Specified by:
      tryPublish in interface StreamBridge
      Parameters:
      event - rule event
      listener - optional listener for evaluation result
      Returns:
      true if published onto the ring buffer, false if saturated
    • evaluateRecord

      public CompletableFuture<StreamResult> evaluateRecord(org.apache.kafka.clients.consumer.ConsumerRecord<String,byte[]> record)
    • stats

      public StreamStats stats()
      Description copied from interface: StreamBridge
      Retrieves current streaming engine metrics and capacity stats.
      Specified by:
      stats in interface StreamBridge
      Returns:
      runtime stats snapshot
    • isRunning

      public boolean isRunning()
      Description copied from interface: StreamBridge
      Checks if the stream bridge is active and accepting events.
      Specified by:
      isRunning in interface StreamBridge
      Returns:
      true if running, false if stopped or shut down
    • shutdown

      public void shutdown()
    • shutdown

      public void shutdown(Duration timeout)
      Description copied from interface: StreamBridge
      Gracefully shuts down the stream bridge within the specified timeout duration.
      Specified by:
      shutdown in interface StreamBridge
      Parameters:
      timeout - maximum time to await in-flight events draining
    • getConsumer

      public KafkaRuleStreamConsumer getConsumer()
    • getProducer

      public KafkaRuleStreamProducer getProducer()