Skip to content

StreamQueueCheckpointer Methods

Coalesces and persists stream queue checkpoints using an IStreamCheckpointStore.

FlushAsync(CancellationToken)

View source
public Task FlushAsync(CancellationToken cancellationToken)
Flushes any pending checkpoint to persistent storage, ensuring the latest offset is durably saved. Called during shutdown or rebalancing to prevent message replay on restart.

Parameters

cancellationTokenCancellationToken
The cancellation token.

Returns

A System.Threading.Tasks.Task representing the flush operation.

Load

View source
[System.Obsolete(Use the overload which accepts a CancellationToken.)]
public Task<string> Load()
Loads the checkpoint.

Returns

The checkpoint.

Load(CancellationToken)

View source
public Task<string> Load(CancellationToken cancellationToken)
Loads the checkpoint.

Parameters

cancellationTokenCancellationToken
The cancellation token.

Returns

The checkpoint.

Update(string, DateTime)

View source
[System.Obsolete(Use the overload which accepts a CancellationToken.)]
public void Update(string offset, DateTime utcNow)
Updates the checkpoint.

Parameters

offsetstring
The offset.
utcNowDateTime
The current UTC time.

Update(string, DateTime, CancellationToken)

View source
public void Update(string offset, DateTime utcNow, CancellationToken cancellationToken)
Updates the checkpoint.

Parameters

offsetstring
The offset.
utcNowDateTime
The current UTC time.
cancellationTokenCancellationToken
The cancellation token.