Class EventsBySliceFirehose$

java.lang.Object
org.apache.pekko.persistence.query.typed.internal.EventsBySliceFirehose$
All Implemented Interfaces:
ExtensionId<org.apache.pekko.persistence.query.typed.internal.EventsBySliceFirehose>, ExtensionIdProvider

public class EventsBySliceFirehose$ extends Object implements ExtensionId<org.apache.pekko.persistence.query.typed.internal.EventsBySliceFirehose>, ExtensionIdProvider
INTERNAL API

The purpose is to share the stream of events from the database and fan out to connected consumer streams. Thereby less queries and loading of events from the database. The shared stream is called the firehose stream.

The fan out of the firehose stream is via a BroadcastHub that consumer streams dynamically attach to.

A new consumer starts a catchup stream since the start offset typically is behind the live firehose stream. In the beginning it will emit events only from the catchup stream. Offset progress for the firehose stream is tracked and when the catchup stream has caught up with the firehose stream it will switch over to emitting from firehose stream and close the catchup stream. During an overlap period of time it will use events from both catchup and firehose streams to make sure that no events are missed. During this overlap time there is best effort deduplication.

The BroadcastHub has a limited buffer that holds events between the slowest and fastest consumer. When the buffer is full the fastest consumer can't progress faster than the slowest. Short periods of slow down can be fine, but after a while the slow consumers are detected and aborted. They have to connect again and try catching up, but without slowing down other streams.

The firehose stream is started on demand when the first consumer is attaching. It will be stopped when the last consumer is stopped, but it stays around for a while to make it more efficient for new or restarted consumers to attach again.

  • Field Details

    • MODULE$

      public static final EventsBySliceFirehose$ MODULE$
      Static reference to the singleton instance of this Scala object.
  • Constructor Details

    • EventsBySliceFirehose$

      public EventsBySliceFirehose$()
  • Method Details

    • get

      public org.apache.pekko.persistence.query.typed.internal.EventsBySliceFirehose get(ActorSystem system)
      Description copied from interface: ExtensionId
      Returns an instance of the extension identified by this ExtensionId instance. Java API For extensions written in Scala that are to be used from Java also, this method should be overridden to get correct return type.
      
       override def get(system: ActorSystem): TheExtension = super.get(system)
       
      Specified by:
      get in interface ExtensionId<org.apache.pekko.persistence.query.typed.internal.EventsBySliceFirehose>
    • get

      public org.apache.pekko.persistence.query.typed.internal.EventsBySliceFirehose get(ClassicActorSystemProvider system)
      Description copied from interface: ExtensionId
      Returns an instance of the extension identified by this ExtensionId instance. Java API For extensions written in Scala that are to be used from Java also, this method should be overridden to get correct return type.
      
       override def get(system: ClassicActorSystemProvider): TheExtension = super.get(system)
       
      Specified by:
      get in interface ExtensionId<org.apache.pekko.persistence.query.typed.internal.EventsBySliceFirehose>
    • lookup

      public EventsBySliceFirehose$ lookup()
      Description copied from interface: ExtensionIdProvider
      Returns the canonical ExtensionId for this Extension
      Specified by:
      lookup in interface ExtensionIdProvider
    • createExtension

      public org.apache.pekko.persistence.query.typed.internal.EventsBySliceFirehose createExtension(ExtendedActorSystem system)
      Description copied from interface: ExtensionId
      Is used by Pekko to instantiate the Extension identified by this ExtensionId, internal use only.
      Specified by:
      createExtension in interface ExtensionId<org.apache.pekko.persistence.query.typed.internal.EventsBySliceFirehose>
    • timestampOffset

      public TimestampOffset timestampOffset(EventEnvelope<Object> env)
    • isDurationGreaterThan

      public boolean isDurationGreaterThan(Instant from, Instant to, Duration duration)