Interface ReaderGroup
-
- All Superinterfaces:
java.lang.AutoCloseable
,ReaderGroupNotificationListener
public interface ReaderGroup extends ReaderGroupNotificationListener, java.lang.AutoCloseable
A reader group is a collection of readers that collectively read all the events in the stream. The events are distributed among the readers in the group such that each event goes to only one reader. The readers in the group may change over time. Readers are added to the group by callingEventStreamClientFactory.createReader(String, String, Serializer, ReaderConfig)
and are removed by callingreaderOffline(String, Position)
-
-
Method Summary
All Methods Instance Methods Abstract Methods Modifier and Type Method Description void
close()
Closes the reader group, freeing any resources associated with it.java.util.concurrent.CompletableFuture<java.util.Map<Stream,StreamCut>>
generateStreamCuts(java.util.concurrent.ScheduledExecutorService backgroundExecutor)
Generates aStreamCut
after co-ordinating with all the readers usingStateSynchronizer
.java.lang.String
getGroupName()
Returns the name of the group.ReaderGroupMetrics
getMetrics()
Returns metrics for this reader group.java.util.Set<java.lang.String>
getOnlineReaders()
Returns a set of readerIds for the readers that are considered to be online by the group.ReaderSegmentDistribution
getReaderSegmentDistribution()
Returns current distribution of number of segments assigned to each reader in the reader group.java.lang.String
getScope()
Returns the scope of the stream which the group is associated with.java.util.Map<Stream,StreamCut>
getStreamCuts()
Returns aStreamCut
for each stream that this reader group is reading from.java.util.Set<java.lang.String>
getStreamNames()
Returns the set of scoped stream names which was used to configure this group.java.util.concurrent.CompletableFuture<Checkpoint>
initiateCheckpoint(java.lang.String checkpointName)
Initiate a checkpoint.java.util.concurrent.CompletableFuture<Checkpoint>
initiateCheckpoint(java.lang.String checkpointName, java.util.concurrent.ScheduledExecutorService backgroundExecutor)
Initiate a checkpoint.void
readerOffline(java.lang.String readerId, Position lastPosition)
Invoked when a reader that was added to the group is no longer consuming events.void
resetReaderGroup()
Reset a reader group to successfully completed last checkpoint.void
resetReaderGroup(ReaderGroupConfig config)
Reset a reader group with the providedReaderGroupConfig
.void
updateRetentionStreamCut(java.util.Map<Stream,StreamCut> streamCuts)
Update Retention Stream-Cut for Streams in this Reader Group.-
Methods inherited from interface io.pravega.client.stream.notifications.ReaderGroupNotificationListener
getEndOfDataNotifier, getSegmentNotifier
-
-
-
-
Method Detail
-
getMetrics
ReaderGroupMetrics getMetrics()
Returns metrics for this reader group.- Returns:
- a ReaderGroupMetrics object for this reader group.
-
getScope
java.lang.String getScope()
Returns the scope of the stream which the group is associated with.- Returns:
- A scope string
-
getGroupName
java.lang.String getGroupName()
Returns the name of the group.- Returns:
- Reader group name
-
initiateCheckpoint
java.util.concurrent.CompletableFuture<Checkpoint> initiateCheckpoint(java.lang.String checkpointName, java.util.concurrent.ScheduledExecutorService backgroundExecutor)
Initiate a checkpoint. This causes all readers in the group to receive a specialEventRead
that contains the provided checkpoint name. This can be used to provide an indication to them that they should persist their state. Once all of the readers have received the notification and resumed reading the future will return aCheckpoint
object which contains the StreamCut of the reader group at the time they received the checkpoint. This can be used to reset the group to this point in the stream by callingresetReaderGroup(ReaderGroupConfig)
if the checkpoint fails or the result cannot be obtained an exception will be set on the future. This method can be called and a new checkpoint can be initiated while another is still in progress if they have different names. If this method is called again before the checkpoint has completed with the same name the future returned to the second caller will refer to the same checkpoint object as the first.- Parameters:
checkpointName
- The name of the checkpoint (For identification purposes)backgroundExecutor
- A threadPool that can be used to poll for the completion of the checkpoint.- Returns:
- A future Checkpoint object that can be used to restore the reader group to this position.
-
initiateCheckpoint
java.util.concurrent.CompletableFuture<Checkpoint> initiateCheckpoint(java.lang.String checkpointName)
Initiate a checkpoint. This causes all readers in the group to receive a specialEventRead
that contains the provided checkpoint name. This can be used to provide an indication to them that they should persist their state. Once all of the readers have received the notification and resumed reading the future will return aCheckpoint
object which contains the StreamCut of the reader group at the time they received the checkpoint. This can be used to reset the group to this point in the stream by callingresetReaderGroup(ReaderGroupConfig)
if the checkpoint fails or the result cannot be obtained an exception will be set on the future. This method can be called and a new checkpoint can be initiated while another is still in progress if they have different names. If this method is called again before the checkpoint has completed with the same name the future returned to the second caller will refer to the same checkpoint object as the first. Client internal thread pool executor is used to poll for the completion of the checkpoint- Parameters:
checkpointName
- The name of the checkpoint (For identification purposes)- Returns:
- A future Checkpoint object that can be used to restore the reader group to this position.
-
resetReaderGroup
void resetReaderGroup(ReaderGroupConfig config)
Reset a reader group with the providedReaderGroupConfig
.- The stream(s) that are part of the reader group can be specified using
ReaderGroupConfig.ReaderGroupConfigBuilder.stream(String)
,ReaderGroupConfig.ReaderGroupConfigBuilder.stream(String, StreamCut)
andReaderGroupConfig.ReaderGroupConfigBuilder.stream(String, StreamCut, StreamCut)
.- To reset a reader group to a given checkpoint use
ReaderGroupConfig.ReaderGroupConfigBuilder.startFromCheckpoint(Checkpoint)
api.- To reset a reader group to a given StreamCut use
All existing readers will have to callReaderGroupConfig.ReaderGroupConfigBuilder.startFromStreamCuts(Map)
.EventStreamClientFactory.createReader(String, String, Serializer, ReaderConfig)
. If they continue to read events they will eventually encounter anReinitializationRequiredException
.- Parameters:
config
- The new configuration for the ReaderGroup. To use a different checkpoint, set the `startingStreamCuts` on the `ReaderGroupConfig` from a streamcut obtained from aCheckpoint
orinitiateCheckpoint(String, ScheduledExecutorService)
.
-
resetReaderGroup
void resetReaderGroup()
Reset a reader group to successfully completed last checkpoint. Successfully completed last checkpoint can be the last checkpoint created when automatic checkpointing is enabled as a part ofReaderGroupConfig
or manually created by callinginitiateCheckpoint(String, ScheduledExecutorService)
If there is no successfully completed Last checkpoint present then this call reset the reader group to the original streamcut from `ReaderGroupConfig`.
-
readerOffline
void readerOffline(java.lang.String readerId, Position lastPosition)
Invoked when a reader that was added to the group is no longer consuming events. This will cause the events that were going to that reader to be redistributed among the other readers. Events after the lastPosition provided will be (re)read by other readers in theReaderGroup
. Note that this method is automatically invoked byEventStreamReader.close()
- Parameters:
readerId
- The id of the reader that is offline.lastPosition
- The position of the last event that was successfully processed by the reader.
-
getOnlineReaders
java.util.Set<java.lang.String> getOnlineReaders()
Returns a set of readerIds for the readers that are considered to be online by the group. i.e.EventStreamClientFactory.createReader(String, String, Serializer, ReaderConfig)
was called butreaderOffline(String, Position)
was not called subsequently.- Returns:
- Set of active reader IDs of the group
-
getStreamNames
java.util.Set<java.lang.String> getStreamNames()
Returns the set of scoped stream names which was used to configure this group.- Returns:
- Set of streams for this group.
-
getStreamCuts
java.util.Map<Stream,StreamCut> getStreamCuts()
Returns aStreamCut
for each stream that this reader group is reading from. The stream cut corresponds to the last checkpointed read offsets of the readers, and it can be used by the application as reference to such a position. A more preciseStreamCut
, with the latest read offsets can be obtained usinggenerateStreamCuts(ScheduledExecutorService)
API.- Returns:
- Map of streams that this group is reading from to the corresponding cuts.
-
generateStreamCuts
@Beta java.util.concurrent.CompletableFuture<java.util.Map<Stream,StreamCut>> generateStreamCuts(java.util.concurrent.ScheduledExecutorService backgroundExecutor)
Generates aStreamCut
after co-ordinating with all the readers usingStateSynchronizer
. AStreamCut
is generated by using the latest segment read offsets returned by the readers along with unassigned segments (if any). The configurationReaderGroupConfig.groupRefreshTimeMillis
decides the maximum delay by which the readers return the latest read offsets of their assigned segments.The
StreamCut
generated by this API can be used by the application as a reference to a position in the stream. This is guaranteed to be greater than or equal to the position of the readers at the point of invocation of the API. TheStreamCut
s generated can be used to perform bounded processing of the Stream by configuring aReaderGroup
with aReaderGroupConfig
where theStreamCut
s are specified as the lower bound and/or upper bounds using the apisReaderGroupConfig.ReaderGroupConfigBuilder.stream(Stream, StreamCut, StreamCut)
orReaderGroupConfig.ReaderGroupConfigBuilder.stream(Stream, StreamCut)
orReaderGroupConfig.ReaderGroupConfigBuilder.startFromStreamCuts(Map)
.Note: Generating a precise
StreamCut
, for example aStreamCut
pointing to end of Q1 across all segments, is difficult as it depends on the configurationReaderGroupConfig.groupRefreshTimeMillis
which decides the duration by which all the readers running on different machines/ processes respond with their latest read offsets. Hence, theStreamCut
would point to a position in theStream
which might include events from Q2. The application thus would need to filter out such additional events.- Parameters:
backgroundExecutor
- A thread pool that will be used to poll if the positions from all the readers have been fetched.- Returns:
- A future to a Map of Streams (that this group is reading from) to its corresponding cuts.
-
updateRetentionStreamCut
void updateRetentionStreamCut(java.util.Map<Stream,StreamCut> streamCuts)
Update Retention Stream-Cut for Streams in this Reader Group. SeeReaderGroupConfig.StreamDataRetention.MANUAL_RELEASE_AT_USER_STREAMCUT
- Parameters:
streamCuts
- A Map with a Stream-Cut for each Stream. StreamCut indicates position in the Stream till which data has been consumed and can be deleted.
-
getReaderSegmentDistribution
ReaderSegmentDistribution getReaderSegmentDistribution()
Returns current distribution of number of segments assigned to each reader in the reader group.- Returns:
- an instance of ReaderSegmentDistribution which describes the distribution of segments to readers including unassigned segments.
-
close
void close()
Closes the reader group, freeing any resources associated with it.- Specified by:
close
in interfacejava.lang.AutoCloseable
-
-