Class KafkaRuleStreamProducer

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

public class KafkaRuleStreamProducer extends Object implements AutoCloseable
High-performance Kafka producer for streaming rule events and structured evaluation results.
  • Constructor Details

    • KafkaRuleStreamProducer

      public KafkaRuleStreamProducer(KafkaStreamConfig config)
    • KafkaRuleStreamProducer

      public KafkaRuleStreamProducer(org.apache.kafka.clients.producer.Producer<String,byte[]> producer)
  • Method Details

    • publish

      public CompletableFuture<org.apache.kafka.clients.producer.RecordMetadata> publish(String topic, String key, byte[] payload)
    • publishResult

      public CompletableFuture<org.apache.kafka.clients.producer.RecordMetadata> publishResult(String outputTopic, StreamResult result)
    • publishEvent

      public CompletableFuture<org.apache.kafka.clients.producer.RecordMetadata> publishEvent(String topic, RuleEvent event)
    • flush

      public void flush()
    • close

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

      public org.apache.kafka.clients.producer.Producer<String,byte[]> getProducer()