Class EventBus
java.lang.Object
com.northconcepts.datapipeline.eventbus.EventBus
An asynchronous, in-memory, event delivery service. When used with
EventBusReader and
EventBusWriter, the bus allows pipelines to be loosely coupled and used in one-to-many (publish-subscribe)
scenarios.- See Also:
-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionstatic enumThe lifecycle states of a bus; listeners, publishers and events are only accepted whileALIVE. -
Field Summary
Fields -
Constructor Summary
ConstructorsConstructorDescriptionEventBus()Creates a bus using a single thread to deliver events.EventBus(ExecutorService executorService) Creates a bus using the specified ExecutorService to deliver events to listeners. -
Method Summary
Modifier and TypeMethodDescriptionstatic voidRegisters a listener notified, via the system event bus, when any event bus is created or begins shutting down.addExceptionListener(EventFilter filter, ExceptionListener listener) Registers an exception listener that only receives the exception events the filter allows.addExceptionListener(ExceptionListener listener) Registers a listener notified when a listener on this bus throws an exception while receiving an event.addListener(EventFilter filter, UntypedEventListener listener) Deprecated.addListener(EventFilter filter, UntypedEventListener listener, Object topic) Deprecated.addListener(UntypedEventListener listener) Deprecated.UseaddUntypedEventListener(UntypedEventListener)instead.<T> EventBusaddListener(Class<T> eventListenerClass, EventFilter filter, T listener) Subscribes a listener to the events of its listener interface that the filter allows (nullallows all).<T> EventBusaddListener(Class<T> eventListenerClass, EventFilter filter, T listener, Object topic) Subscribes a listener to the events of its listener interface that the filter allows and, unlesstopicis null, that carry an equal topic.<T> EventBusaddListener(Class<T> eventListenerClass, T listener) Subscribes a listener to every event of its listener interface, whatever the topic.addShutdownListener(EventFilter filter, ShutdownListener listener) Registers a listener notified when this bus begins shutting down; the filter may benull.addUntypedEventListener(EventFilter filter, UntypedEventListener listener) Subscribes a listener to every event on this bus that the filter allows (nullallows all).addUntypedEventListener(EventFilter filter, UntypedEventListener listener, Object topic) Subscribes a listener to every event on this bus that the filter allows and, unlesstopicis null, that carries an equal topic.addUntypedEventListener(UntypedEventListener listener) Subscribes a listener to every event published on this bus, whatever its listener interface or topic.voidAsserts that this endpoint has started reading or writing, otherwise an exception is thrown.protected voidRemoves the listeners whose soft references were garbage collected; called after each delivery.protected Objectclone()protected <T> voiddeliverEvent(Event<T> event) Delivers the event to the matching typed listeners, then to the untyped listeners, withgetCurrentEvent()set; runs on an executor thread.protected voidfinalize()getAll()Returns all the event buses in the JVM/classloader that have not been terminated.longReturns the approximate total number of tasks that have completed execution by this bus' executorService or -1 if unknown.longReturns the core number of threads in this bus' executorService or -1 if unknown.The time this job instance was created.static Event<?> Returns the event currently being delivered on this thread.longReturns the number of exceptions this bus has experienced.longReturns the approximate number of threads that are actively delivering events using this bus' executorService or -1 if unknown.longReturns the combined number of events each listener on this event bus has received (typicallygetEventsPublished()*getListenerCount()).longReturns the number of events sent to this event bus for delivery to listeners.longReturns the number of events waiting to be executed by this bus' executorService or -1 if unknown.longgetId()Returns the unique (sequential) ID of this event bus (within this JVM/classloader).longReturns the largest number of threads that have ever simultaneously been in the pool in this bus' executorService or -1 if unknown.intReturns the total number of listeners subscribed to this bus.longReturns the maximum allowed number of threads in this bus' executorService or -1 if unknown.getName()Returns the name assigned to this bus or its default name if one was not assigned.longReturns the current number of threads in this bus' executorService pool or -1 if unknown.<T> TgetPublisher(Class<T> eventListenerClass) Returns a proxy of the listener interface that publishes each call as an event on this bus, without a source or topic.<T> TgetPublisher(Object eventSource, Class<T> eventListenerClass) Returns a proxy of the listener interface that publishes each call as an event on this bus, from the given source and without a topic.<T> TgetPublisher(Object eventSource, Class<T> eventListenerClass, Object topic) Returns a proxy of the listener interface that publishes each call as an event on this bus, with the given source and topic; the bus's executor delivers the events to the listeners.longReturns the number of event publishers created by this bus.longReturns the time in milliseconds this bus will wait for events to complete before forcibly shutting down (defaults to 30 seconds).Returns time as human readable string this bus will wait for events to complete before forcibly shutting down (defaults to 30 seconds).getState()Returns the current state of this bus.static EventBusReturns the shared JVM-wide bus used for lifecycle events; a JVM shutdown hook shuts it down.longReturns the approximate total number of tasks that have ever been scheduled for execution by this bus' executorService or -1 if unknown.intReturns the number of type-specific listeners subscribed to this bus.protected <T> TypedListenerList<T> getTypedListenerList(Class<T> eventListenerClass) Returns the listeners registered for the given listener interface ornullif none has been added for it.intReturns the number of non-type-specific listeners subscribed to this bus.intDeprecated.UsegetUntypedEventListenerCount()instead.getUuid()Returns the universally unique identifier (UUID) for this job.protected voidhandleException(Event<?> event, Object listener, Throwable e) Counts a listener's exception and publishes it to the exception listeners, unless an exception listener threw it.protected longCounts one more delivery to a listener and returns the newgetEventsDelivered()total.booleanisAlive()Returns true if this bus has not been shutdown.booleanisDebug()Indicates if the bus should emit debug log messages when events are published or delivered and during the shutdown process.booleanReturns true if this bus has been shutdown or is in the process of shutting down.booleanIndicates if the bus should publish its own EventBusLifecycleListener and ShutdownListener events (defaults totrue).booleanIndicates if the event bus uses strong references (default) or soft references for holding listeners added to the bus.protected <T> voidpublishEvent(Event<T> event) Numbers the event and hands it to the executor for delivery; throws if this bus is not alive, except for lifecycle and shutdown events, which are dropped.static voidremoveExceptionListener(ExceptionListener listener) removeListener(UntypedEventListener listener) Deprecated.UseremoveUntypedEventListener(UntypedEventListener)instead.<T> EventBusremoveListener(Class<T> eventListenerClass, T listener) removeListener(Object listener) Unsubscribes the listener from every listener interface it was added for; untyped listeners are not affected.<T> EventBusremoveListenersByTopic(Class<T> eventListenerClass, Object topic) Unsubscribes the given interface's listeners that were added with a topic equal to this one.removeShutdownListener(ShutdownListener listener) protected EventBusremoveTypedListenerReference(Reference<?> reference) Removes the typed listener held by the given reference once the garbage collector has cleared it.protected EventBusremoveUntypedEventListenerReference(Reference<?> reference) Removes the untyped listener held by the given reference once the garbage collector has cleared it.Unsubscribes the untyped listeners that were added with a topic equal to this one.protected EventBusremoveUntypedListenerReference(Reference<?> reference) Deprecated.UseremoveUntypedEventListenerReference(Reference)instead.setDebug(boolean debug) Indicates if the bus should emit debug log messages when events are published or delivered and during the shutdown process.Assigns a new name to this bus.setPublishLifecycleEvents(boolean publishLifecycleEvents) Indicates if the bus should publish its own EventBusLifecycleListener and ShutdownListener events (defaults totrue).setShutdownTimeout(long shutdownTimeout) Sets the time in milliseconds this bus will wait for events to complete before forcibly shutting down (defaults to 30 seconds).setUseStrongListenerReference(boolean useStrongListenerReference) Indicates if the event bus uses strong references (default) or soft references for holding listeners added to the bus.voidshutdown()Notifies the shutdown and lifecycle listeners, waits up togetShutdownTimeout()milliseconds for queued events, then stops the executor and removes all listeners.
-
Field Details
-
log
public static final Logger log
-
-
Constructor Details
-
EventBus
Creates a bus using the specified ExecutorService to deliver events to listeners. -
EventBus
public EventBus()Creates a bus using a single thread to deliver events. This configuration allows the bus to deliver events sequentially to listeners, in the order they were published, while not blocking the publisher.
-
-
Method Details
-
getAll
Returns all the event buses in the JVM/classloader that have not been terminated. -
getSystemEventBus
Returns the shared JVM-wide bus used for lifecycle events; a JVM shutdown hook shuts it down. -
addEventBusLifecycleListener
Registers a listener notified, via the system event bus, when any event bus is created or begins shutting down. -
removeEventBusLifecycleListener
-
getCurrentEvent
Returns the event currently being delivered on this thread. -
getId
public long getId()Returns the unique (sequential) ID of this event bus (within this JVM/classloader). -
getUuid
Returns the universally unique identifier (UUID) for this job. -
getName
Returns the name assigned to this bus or its default name if one was not assigned. -
setName
Assigns a new name to this bus. -
getState
Returns the current state of this bus. -
isAlive
public boolean isAlive()Returns true if this bus has not been shutdown. This is equivalent togetState()==EventBus.State.ALIVE -
isNotAlive
public boolean isNotAlive()Returns true if this bus has been shutdown or is in the process of shutting down. This is equivalent togetState()!=EventBus.State.ALIVE -
assertAlive
public void assertAlive()Asserts that this endpoint has started reading or writing, otherwise an exception is thrown. -
getCreatedOn
The time this job instance was created. -
isDebug
public boolean isDebug()Indicates if the bus should emit debug log messages when events are published or delivered and during the shutdown process. -
setDebug
Indicates if the bus should emit debug log messages when events are published or delivered and during the shutdown process. -
getShutdownTimeout
public long getShutdownTimeout()Returns the time in milliseconds this bus will wait for events to complete before forcibly shutting down (defaults to 30 seconds). -
getShutdownTimeoutAsString
Returns time as human readable string this bus will wait for events to complete before forcibly shutting down (defaults to 30 seconds). -
isPublishLifecycleEvents
public boolean isPublishLifecycleEvents()Indicates if the bus should publish its own EventBusLifecycleListener and ShutdownListener events (defaults totrue). -
setPublishLifecycleEvents
Indicates if the bus should publish its own EventBusLifecycleListener and ShutdownListener events (defaults totrue). -
setShutdownTimeout
Sets the time in milliseconds this bus will wait for events to complete before forcibly shutting down (defaults to 30 seconds). -
isUseStrongListenerReference
public boolean isUseStrongListenerReference()Indicates if the event bus uses strong references (default) or soft references for holding listeners added to the bus. Soft references are automatically garbage collected when the JVM runs out of memory in an effort to free up resources. -
setUseStrongListenerReference
Indicates if the event bus uses strong references (default) or soft references for holding listeners added to the bus. Soft references are automatically garbage collected when the JVM runs out of memory in an effort to free up resources. -
getPublishersCreated
public long getPublishersCreated()Returns the number of event publishers created by this bus. -
getEventsPublished
public long getEventsPublished()Returns the number of events sent to this event bus for delivery to listeners. -
getEventsDelivered
public long getEventsDelivered()Returns the combined number of events each listener on this event bus has received (typicallygetEventsPublished()*getListenerCount()). -
incEventsDelivered
protected long incEventsDelivered()Counts one more delivery to a listener and returns the newgetEventsDelivered()total. -
getEventsActive
public long getEventsActive()Returns the approximate number of threads that are actively delivering events using this bus' executorService or -1 if unknown. -
getEventsQueued
public long getEventsQueued()Returns the number of events waiting to be executed by this bus' executorService or -1 if unknown. -
getTaskCount
public long getTaskCount()Returns the approximate total number of tasks that have ever been scheduled for execution by this bus' executorService or -1 if unknown. -
getCompletedTaskCount
public long getCompletedTaskCount()Returns the approximate total number of tasks that have completed execution by this bus' executorService or -1 if unknown. -
getCorePoolSize
public long getCorePoolSize()Returns the core number of threads in this bus' executorService or -1 if unknown. -
getLargestPoolSize
public long getLargestPoolSize()Returns the largest number of threads that have ever simultaneously been in the pool in this bus' executorService or -1 if unknown. -
getMaximumPoolSize
public long getMaximumPoolSize()Returns the maximum allowed number of threads in this bus' executorService or -1 if unknown. -
getPoolSize
public long getPoolSize()Returns the current number of threads in this bus' executorService pool or -1 if unknown. -
getErrorCount
public long getErrorCount()Returns the number of exceptions this bus has experienced. -
addExceptionListener
Registers a listener notified when a listener on this bus throws an exception while receiving an event. -
addExceptionListener
Registers an exception listener that only receives the exception events the filter allows. -
removeExceptionListener
-
addShutdownListener
Registers a listener notified when this bus begins shutting down; the filter may benull. -
removeShutdownListener
-
getTypedListenerList
Returns the listeners registered for the given listener interface ornullif none has been added for it. -
addListener
Subscribes a listener to every event of its listener interface, whatever the topic. -
addListener
Subscribes a listener to the events of its listener interface that the filter allows (nullallows all). -
addListener
public <T> EventBus addListener(Class<T> eventListenerClass, EventFilter filter, T listener, Object topic) Subscribes a listener to the events of its listener interface that the filter allows and, unlesstopicis null, that carry an equal topic. -
removeListener
-
removeListenersByTopic
Unsubscribes the given interface's listeners that were added with a topic equal to this one. -
removeListener
Unsubscribes the listener from every listener interface it was added for; untyped listeners are not affected. -
removeTypedListenerReference
Removes the typed listener held by the given reference once the garbage collector has cleared it. -
getTypedListenerCount
public int getTypedListenerCount()Returns the number of type-specific listeners subscribed to this bus. -
addListener
Deprecated.UseaddUntypedEventListener(UntypedEventListener)instead.- Returns:
-
addUntypedEventListener
Subscribes a listener to every event published on this bus, whatever its listener interface or topic. -
addListener
Deprecated. -
addUntypedEventListener
Subscribes a listener to every event on this bus that the filter allows (nullallows all). -
addListener
@Deprecated public EventBus addListener(EventFilter filter, UntypedEventListener listener, Object topic) Deprecated. -
addUntypedEventListener
public EventBus addUntypedEventListener(EventFilter filter, UntypedEventListener listener, Object topic) Subscribes a listener to every event on this bus that the filter allows and, unlesstopicis null, that carries an equal topic. -
removeListener
Deprecated.UseremoveUntypedEventListener(UntypedEventListener)instead. -
removeUntypedEventListener
-
removeUntypedEventListenersByTopic
Unsubscribes the untyped listeners that were added with a topic equal to this one. -
removeUntypedListenerReference
Deprecated.UseremoveUntypedEventListenerReference(Reference)instead. -
removeUntypedEventListenerReference
Removes the untyped listener held by the given reference once the garbage collector has cleared it. -
getUntypedListenerCount
Deprecated.UsegetUntypedEventListenerCount()instead.Returns the number of non-type-specific listeners subscribed to this bus. -
getUntypedEventListenerCount
public int getUntypedEventListenerCount()Returns the number of non-type-specific listeners subscribed to this bus. -
getListenerCount
public int getListenerCount()Returns the total number of listeners subscribed to this bus. -
getPublisher
Returns a proxy of the listener interface that publishes each call as an event on this bus, without a source or topic. -
getPublisher
Returns a proxy of the listener interface that publishes each call as an event on this bus, from the given source and without a topic. -
getPublisher
Returns a proxy of the listener interface that publishes each call as an event on this bus, with the given source and topic; the bus's executor delivers the events to the listeners. -
publishEvent
Numbers the event and hands it to the executor for delivery; throws if this bus is not alive, except for lifecycle and shutdown events, which are dropped. -
deliverEvent
Delivers the event to the matching typed listeners, then to the untyped listeners, withgetCurrentEvent()set; runs on an executor thread.- Throws:
Throwable
-
cleanupGarbage
protected void cleanupGarbage()Removes the listeners whose soft references were garbage collected; called after each delivery. -
handleException
Counts a listener's exception and publishes it to the exception listeners, unless an exception listener threw it. -
finalize
-
shutdown
Notifies the shutdown and lifecycle listeners, waits up togetShutdownTimeout()milliseconds for queued events, then stops the executor and removes all listeners. Does nothing if this bus is not alive.- Throws:
InterruptedException
-
clone
- Overrides:
clonein classObject- Throws:
CloneNotSupportedException
-
addUntypedEventListener(EventFilter, UntypedEventListener)instead.