Package com.helix.core.stream.kafka
Class KafkaStreamEngine
java.lang.Object
com.helix.core.stream.kafka.KafkaStreamEngine
- All Implemented Interfaces:
StreamBridge,AutoCloseable
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 Summary
ConstructorsConstructorDescriptionKafkaStreamEngine(KafkaStreamConfig config) KafkaStreamEngine(KafkaStreamConfig config, org.apache.kafka.clients.consumer.Consumer<String, byte[]> consumer, org.apache.kafka.clients.producer.Producer<String, byte[]> producer) -
Method Summary
Modifier and TypeMethodDescriptionevaluateRecord(org.apache.kafka.clients.consumer.ConsumerRecord<String, byte[]> record) booleanChecks if the stream bridge is active and accepting events.Publishes a rule evaluation event into the stream and returns a CompletableFuture for the result.publish(StreamRecord<T> record) Publishes a generic stream record into the stream and returns a CompletableFuture for the result.voidregisterRule(String topicOrRuleName, CompiledRule rule) voidshutdown()voidGracefully shuts down the stream bridge within the specified timeout duration.voidstart()stats()Retrieves current streaming engine metrics and capacity stats.voidsubscribe(String topic, RuleStreamHandler handler) Subscribes a functional rule stream handler to a given topic.voidsubscribeListener(String topic, StreamListener listener) Subscribes an asynchronous result listener to a given topic.booleantryPublish(RuleEvent event, StreamListener listener) Fast-path non-blocking publish of a rule event with a completion listener, avoiding Future allocation.Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitMethods inherited from interface com.helix.api.stream.StreamBridge
addListener, close, tryPublish
-
Constructor Details
-
KafkaStreamEngine
-
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
Description copied from interface:StreamBridgeSubscribes a functional rule stream handler to a given topic.- Specified by:
subscribein interfaceStreamBridge- Parameters:
topic- stream topichandler- handler invoked for incoming events on the topic
-
subscribeListener
Description copied from interface:StreamBridgeSubscribes an asynchronous result listener to a given topic.- Specified by:
subscribeListenerin interfaceStreamBridge- Parameters:
topic- stream topiclistener- listener invoked when events on the topic complete evaluation
-
registerRule
-
publish
Description copied from interface:StreamBridgePublishes a rule evaluation event into the stream and returns a CompletableFuture for the result.- Specified by:
publishin interfaceStreamBridge- Parameters:
event- the rule event to publish- Returns:
- CompletableFuture completing with the evaluation result
-
publish
Description copied from interface:StreamBridgePublishes a generic stream record into the stream and returns a CompletableFuture for the result.- Specified by:
publishin interfaceStreamBridge- Type Parameters:
T- payload type- Parameters:
record- the stream record to publish- Returns:
- CompletableFuture completing with the evaluation result
-
tryPublish
Description copied from interface:StreamBridgeFast-path non-blocking publish of a rule event with a completion listener, avoiding Future allocation.- Specified by:
tryPublishin interfaceStreamBridge- Parameters:
event- rule eventlistener- 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
Description copied from interface:StreamBridgeRetrieves current streaming engine metrics and capacity stats.- Specified by:
statsin interfaceStreamBridge- Returns:
- runtime stats snapshot
-
isRunning
public boolean isRunning()Description copied from interface:StreamBridgeChecks if the stream bridge is active and accepting events.- Specified by:
isRunningin interfaceStreamBridge- Returns:
- true if running, false if stopped or shut down
-
shutdown
public void shutdown() -
shutdown
Description copied from interface:StreamBridgeGracefully shuts down the stream bridge within the specified timeout duration.- Specified by:
shutdownin interfaceStreamBridge- Parameters:
timeout- maximum time to await in-flight events draining
-
getConsumer
-
getProducer
-