Package com.helix.core.stream.disruptor
Class DisruptorStreamBridge
java.lang.Object
com.helix.core.stream.disruptor.DisruptorStreamBridge
- All Implemented Interfaces:
StreamBridge,AutoCloseable
Ultra-low latency, zero-allocation in-process streaming engine backed by the
LMAX Disruptor ring buffer pattern with cache-line padded events and virtual thread
event processing.
-
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptioncom.lmax.disruptor.RingBuffer<RuleEventHolder> 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.voidpublish(String topic, String ruleName, ExecutionContext context, StreamListener listener) Steady-state zero-allocation publication without intermediate RuleEvent heap allocations.voidregisterRule(String topicOrRuleName, CompiledRule rule) Registers a pre-compiled rule directly to evaluate events for a given topic or rule name.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
-
DisruptorStreamBridge
public DisruptorStreamBridge() -
DisruptorStreamBridge
-
-
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
Registers a pre-compiled rule directly to evaluate events for a given topic or rule name.- Parameters:
topicOrRuleName- topic or rule name identifierrule- pre-compiled rule instance
-
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
-
publish
public void publish(String topic, String ruleName, ExecutionContext context, StreamListener listener) Steady-state zero-allocation publication without intermediate RuleEvent heap allocations.- Parameters:
topic- stream topicruleName- rule namecontext- execution contextlistener- optional result callback listener
-
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
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
-
getRingBuffer
-