Package com.helix.api.stream
Interface StreamBridge
- All Superinterfaces:
AutoCloseable
- All Known Implementing Classes:
DisruptorStreamBridge,KafkaStreamEngine
High-performance streaming SPI bridge for in-process event ingestion,
rule evaluation, and result dispatching.
-
Method Summary
Modifier and TypeMethodDescriptiondefault voidaddListener(String topic, StreamListener listener) Alias forsubscribeListener(String, StreamListener).default voidclose()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.voidGracefully shuts down the stream bridge within the specified timeout duration.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.default booleantryPublish(RuleEvent event) Fast-path fire-and-forget publish of an event.booleantryPublish(RuleEvent event, StreamListener listener) Fast-path non-blocking publish of a rule event with a completion listener, avoiding Future allocation.
-
Method Details
-
subscribe
Subscribes a functional rule stream handler to a given topic.- Parameters:
topic- stream topichandler- handler invoked for incoming events on the topic
-
subscribeListener
Subscribes an asynchronous result listener to a given topic.- Parameters:
topic- stream topiclistener- listener invoked when events on the topic complete evaluation
-
addListener
Alias forsubscribeListener(String, StreamListener).- Parameters:
topic- stream topiclistener- listener invoked when events on the topic complete evaluation
-
publish
Publishes a rule evaluation event into the stream and returns a CompletableFuture for the result.- Parameters:
event- the rule event to publish- Returns:
- CompletableFuture completing with the evaluation result
-
publish
Publishes a generic stream record into the stream and returns a CompletableFuture for the result.- Type Parameters:
T- payload type- Parameters:
record- the stream record to publish- Returns:
- CompletableFuture completing with the evaluation result
-
tryPublish
Fast-path non-blocking publish of a rule event with a completion listener, avoiding Future allocation.- Parameters:
event- rule eventlistener- optional listener for evaluation result- Returns:
- true if published onto the ring buffer, false if saturated
-
tryPublish
Fast-path fire-and-forget publish of an event.- Parameters:
event- rule event- Returns:
- true if accepted, false if saturated
-
stats
StreamStats stats()Retrieves current streaming engine metrics and capacity stats.- Returns:
- runtime stats snapshot
-
isRunning
boolean isRunning()Checks if the stream bridge is active and accepting events.- Returns:
- true if running, false if stopped or shut down
-
shutdown
Gracefully shuts down the stream bridge within the specified timeout duration.- Parameters:
timeout- maximum time to await in-flight events draining
-
close
default void close()- Specified by:
closein interfaceAutoCloseable
-