Class InMemoryEventStorageEngine

java.lang.Object
org.axonframework.eventsourcing.eventstore.inmemory.InMemoryEventStorageEngine
All Implemented Interfaces:
EventStorageEngine

public class InMemoryEventStorageEngine extends Object implements EventStorageEngine
Thread-safe event storage engine that stores events and snapshots in memory.
Since:
3.0
Author:
Rene de Waele
  • Constructor Details

    • InMemoryEventStorageEngine

      public InMemoryEventStorageEngine()
      Initializes an InMemoryEventStorageEngine. The engine will be empty, and there is no offset for the first token.
    • InMemoryEventStorageEngine

      public InMemoryEventStorageEngine(long offset)
      Initializes an InMemoryEventStorageEngine using given offset to initialize the tokens with.
      Parameters:
      offset - The value to use for the token of the first event appended
  • Method Details

    • appendEvents

      public void appendEvents(@Nonnull List<? extends EventMessage<?>> events)
      Description copied from interface: EventStorageEngine
      Append a list of events to the event storage. Events will be appended in the order that they are offered in.

      Note that all events should have a unique event identifier. When storing domain events events should also have a unique combination of aggregate id and sequence number.

      Specified by:
      appendEvents in interface EventStorageEngine
      Parameters:
      events - Events to append to the event storage
    • storeSnapshot

      public void storeSnapshot(@Nonnull DomainEventMessage<?> snapshot)
      Description copied from interface: EventStorageEngine
      Store an event that contains a snapshot of an aggregate. If the event storage already contains a snapshot for the same aggregate, then it will be replaced with the given snapshot.
      Specified by:
      storeSnapshot in interface EventStorageEngine
      Parameters:
      snapshot - The snapshot event of the aggregate that is to be stored
    • readEvents

      public Stream<? extends TrackedEventMessage<?>> readEvents(TrackingToken trackingToken, boolean mayBlock)
      Open an event stream containing all events stored since given tracking token. The returned stream is comprised of events from aggregates as well as other application events. Pass a trackingToken of null to open a stream containing all available events.

      If the value of the given mayBlock is true the returned stream is allowed to block while waiting for new event messages if the end of the stream is reached.

      This implementation produces non-blocking event streams.

      Specified by:
      readEvents in interface EventStorageEngine
      Parameters:
      trackingToken - Object describing the global index of the last processed event or null to create a stream of all events in the store
      mayBlock - If true the storage engine may optionally choose to block to wait for new event messages if the end of the stream is reached.
      Returns:
      A stream containing all tracked event messages stored since the given tracking token
    • readEvents

      public DomainEventStream readEvents(@Nonnull String aggregateIdentifier, long firstSequenceNumber)
      Description copied from interface: EventStorageEngine
      Get a DomainEventStream containing all events published by the aggregate with given aggregateIdentifier starting with the first event having a sequence number that is equal or larger than the given firstSequenceNumber.

      The returned stream is finite, i.e. it should not block to wait for further events if the end of the event stream of the aggregate is reached.

      Specified by:
      readEvents in interface EventStorageEngine
      Parameters:
      aggregateIdentifier - The identifier of the aggregate
      firstSequenceNumber - The expected sequence number of the first event in the returned stream
      Returns:
      A non-blocking DomainEventStream of the given aggregate
    • readSnapshot

      public Optional<DomainEventMessage<?>> readSnapshot(@Nonnull String aggregateIdentifier)
      Description copied from interface: EventStorageEngine
      Try to load a snapshot event of the aggregate with given aggregateIdentifier. If the storage engine has no snapshot event of the aggregate, an empty Optional is returned.
      Specified by:
      readSnapshot in interface EventStorageEngine
      Parameters:
      aggregateIdentifier - The identifier of the aggregate
      Returns:
      An optional with a snapshot of the aggregate
    • createTailToken

      public TrackingToken createTailToken()
      Description copied from interface: EventStorageEngine
      Creates a token that is at the tail of an event stream - that tracks events from the beginning of time.
      Specified by:
      createTailToken in interface EventStorageEngine
      Returns:
      a tracking token at the tail of an event stream, if event stream is empty null is returned
    • createHeadToken

      public TrackingToken createHeadToken()
      Description copied from interface: EventStorageEngine
      Creates a token that is at the head of an event stream - that tracks all new events.
      Specified by:
      createHeadToken in interface EventStorageEngine
      Returns:
      a tracking token at the head of an event stream, if event stream is empty null is returned
    • createTokenAt

      public TrackingToken createTokenAt(@Nonnull Instant dateTime)
      Description copied from interface: EventStorageEngine
      Creates a token that tracks all events after given dateTime. If there is an event exactly at the given dateTime, it will be tracked too.
      Specified by:
      createTokenAt in interface EventStorageEngine
      Parameters:
      dateTime - The date and time for determining criteria how the tracking token should be created. A tracking token should point to very first event before this date and time.
      Returns:
      a tracking token at the given dateTime, if there aren't events matching this criteria null is returned
    • nextTrackingToken

      protected GlobalSequenceTrackingToken nextTrackingToken()
      Returns the tracking token to use for the next event to be stored.
      Returns:
      the tracking token for the next event