Package com.helix.core.stream.kafka
Class KafkaRuleStreamProducer
java.lang.Object
com.helix.core.stream.kafka.KafkaRuleStreamProducer
- All Implemented Interfaces:
AutoCloseable
High-performance Kafka producer for streaming rule events and structured evaluation results.
-
Constructor Summary
ConstructorsConstructorDescriptionKafkaRuleStreamProducer(org.apache.kafka.clients.producer.Producer<String, byte[]> producer) -
Method Summary
Modifier and TypeMethodDescriptionvoidclose()voidflush()org.apache.kafka.clients.producer.Producer<String, byte[]> CompletableFuture<org.apache.kafka.clients.producer.RecordMetadata> CompletableFuture<org.apache.kafka.clients.producer.RecordMetadata> publishEvent(String topic, RuleEvent event) CompletableFuture<org.apache.kafka.clients.producer.RecordMetadata> publishResult(String outputTopic, StreamResult result)
-
Constructor Details
-
KafkaRuleStreamProducer
-
KafkaRuleStreamProducer
-
-
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:
closein interfaceAutoCloseable
-
getProducer
-