Class EventBus

java.lang.Object
com.northconcepts.datapipeline.eventbus.EventBus

public class EventBus extends Object
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:
  • Field Details

    • log

      public static final Logger log
  • Constructor Details

    • EventBus

      public EventBus(ExecutorService executorService)
      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

      public static List<EventBus> getAll()
      Returns all the event buses in the JVM/classloader that have not been terminated.
    • getSystemEventBus

      public static EventBus getSystemEventBus()
      Returns the shared JVM-wide bus used for lifecycle events; a JVM shutdown hook shuts it down.
    • addEventBusLifecycleListener

      public static void addEventBusLifecycleListener(EventBusLifecycleListener listener)
      Registers a listener notified, via the system event bus, when any event bus is created or begins shutting down.
    • removeEventBusLifecycleListener

      public static void removeEventBusLifecycleListener(EventBusLifecycleListener listener)
    • getCurrentEvent

      public static Event<?> 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

      public UUID getUuid()
      Returns the universally unique identifier (UUID) for this job.
    • getName

      public String getName()
      Returns the name assigned to this bus or its default name if one was not assigned.
    • setName

      public EventBus setName(String name)
      Assigns a new name to this bus.
    • getState

      public EventBus.State getState()
      Returns the current state of this bus.
    • isAlive

      public boolean isAlive()
      Returns true if this bus has not been shutdown. This is equivalent to getState() == 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 to getState() != EventBus.State.ALIVE
    • assertAlive

      public void assertAlive()
      Asserts that this endpoint has started reading or writing, otherwise an exception is thrown.
    • getCreatedOn

      public Date 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

      public EventBus setDebug(boolean debug)
      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

      public String 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 to true).
    • setPublishLifecycleEvents

      public EventBus setPublishLifecycleEvents(boolean publishLifecycleEvents)
      Indicates if the bus should publish its own EventBusLifecycleListener and ShutdownListener events (defaults to true).
    • setShutdownTimeout

      public EventBus setShutdownTimeout(long shutdownTimeout)
      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

      public EventBus setUseStrongListenerReference(boolean useStrongListenerReference)
      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 (typically getEventsPublished() * getListenerCount()).
    • incEventsDelivered

      protected long incEventsDelivered()
      Counts one more delivery to a listener and returns the new getEventsDelivered() 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

      public EventBus addExceptionListener(ExceptionListener listener)
      Registers a listener notified when a listener on this bus throws an exception while receiving an event.
    • addExceptionListener

      public EventBus addExceptionListener(EventFilter filter, ExceptionListener listener)
      Registers an exception listener that only receives the exception events the filter allows.
    • removeExceptionListener

      public EventBus removeExceptionListener(ExceptionListener listener)
    • addShutdownListener

      public EventBus addShutdownListener(EventFilter filter, ShutdownListener listener)
      Registers a listener notified when this bus begins shutting down; the filter may be null.
    • removeShutdownListener

      public EventBus removeShutdownListener(ShutdownListener listener)
    • getTypedListenerList

      protected <T> TypedListenerList<T> getTypedListenerList(Class<T> eventListenerClass)
      Returns the listeners registered for the given listener interface or null if none has been added for it.
    • addListener

      public <T> EventBus addListener(Class<T> eventListenerClass, T listener)
      Subscribes a listener to every event of its listener interface, whatever the topic.
    • addListener

      public <T> EventBus addListener(Class<T> eventListenerClass, EventFilter filter, T listener)
      Subscribes a listener to the events of its listener interface that the filter allows (null allows 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, unless topic is null, that carry an equal topic.
    • removeListener

      public <T> EventBus removeListener(Class<T> eventListenerClass, T listener)
    • removeListenersByTopic

      public <T> EventBus removeListenersByTopic(Class<T> eventListenerClass, Object topic)
      Unsubscribes the given interface's listeners that were added with a topic equal to this one.
    • removeListener

      public EventBus removeListener(Object listener)
      Unsubscribes the listener from every listener interface it was added for; untyped listeners are not affected.
    • removeTypedListenerReference

      protected EventBus removeTypedListenerReference(Reference<?> reference)
      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 public EventBus addListener(UntypedEventListener listener)
      Returns:
    • addUntypedEventListener

      public EventBus addUntypedEventListener(UntypedEventListener listener)
      Subscribes a listener to every event published on this bus, whatever its listener interface or topic.
    • addListener

      @Deprecated public EventBus addListener(EventFilter filter, UntypedEventListener listener)
    • addUntypedEventListener

      public EventBus addUntypedEventListener(EventFilter filter, UntypedEventListener listener)
      Subscribes a listener to every event on this bus that the filter allows (null allows all).
    • addListener

      @Deprecated public EventBus addListener(EventFilter filter, UntypedEventListener listener, Object topic)
    • addUntypedEventListener

      public EventBus addUntypedEventListener(EventFilter filter, UntypedEventListener listener, Object topic)
      Subscribes a listener to every event on this bus that the filter allows and, unless topic is null, that carries an equal topic.
    • removeListener

      @Deprecated public EventBus removeListener(UntypedEventListener listener)
    • removeUntypedEventListener

      public EventBus removeUntypedEventListener(UntypedEventListener listener)
    • removeUntypedEventListenersByTopic

      public EventBus removeUntypedEventListenersByTopic(Object topic)
      Unsubscribes the untyped listeners that were added with a topic equal to this one.
    • removeUntypedListenerReference

      @Deprecated protected EventBus removeUntypedListenerReference(Reference<?> reference)
    • removeUntypedEventListenerReference

      protected EventBus removeUntypedEventListenerReference(Reference<?> reference)
      Removes the untyped listener held by the given reference once the garbage collector has cleared it.
    • getUntypedListenerCount

      @Deprecated public int getUntypedListenerCount()
      Deprecated.
      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

      public <T> T getPublisher(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.
    • getPublisher

      public <T> T getPublisher(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.
    • getPublisher

      public <T> T getPublisher(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.
    • publishEvent

      protected <T> void publishEvent(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.
    • deliverEvent

      protected <T> void deliverEvent(Event<T> event) throws Throwable
      Delivers the event to the matching typed listeners, then to the untyped listeners, with getCurrentEvent() 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

      protected void handleException(Event<?> event, Object listener, Throwable e)
      Counts a listener's exception and publishes it to the exception listeners, unless an exception listener threw it.
    • finalize

      protected void finalize() throws Throwable
      Overrides:
      finalize in class Object
      Throws:
      Throwable
    • shutdown

      public void shutdown() throws InterruptedException
      Notifies the shutdown and lifecycle listeners, waits up to getShutdownTimeout() milliseconds for queued events, then stops the executor and removes all listeners. Does nothing if this bus is not alive.
      Throws:
      InterruptedException
    • clone

      protected Object clone() throws CloneNotSupportedException
      Overrides:
      clone in class Object
      Throws:
      CloneNotSupportedException