BufferPullConsumer<TItem> is the one shared pull-consumer implementation. It is
not storage-specific. It:
- keeps the assigned partition list;
- selects partitions in round-robin order;
- pulls batches from partitions;
- tries other assigned partitions when the selected partition has no data;
- waits asynchronously when no assigned partition has data;
- commits manually consumed batches;
- auto-commits when
AutoCommitis enabled.
The consumer scans assigned partitions in a ring. After a successful pull, the next scan starts after the partition that supplied the batch. Empty scans leave this cursor unchanged. Continuously readable partitions receive one pull opportunity per rotation, including when other partitions are empty. This is batch scheduling fairness, not equal processing time, per-key fairness, or global ordering. Manual commits still apply to the most recently delivered batch; an uncommitted partition can replay its batch when its next turn arrives. The cursor stays within the assigned array bounds.
Consumption returns an asynchronous stream of batches:
await foreach (var batch in consumer.ConsumeAsync(cancellationToken))
{
foreach (var item in batch)
{
// Process item.
}
await consumer.CommitAsync();
}When AutoCommit is false, the consumer read position does not advance until
CommitAsync is called. This gives at-least-once behavior within the current process.
In MemoryMappedFile mode it also applies across process restarts for records that have reached a
flush boundary; CommitAsync itself forces that boundary.
The relationship between groups, consumers, partitions, ordering, and assignment is described in Partitioning and concurrency.
Consumers do not spin when no data is available. The common consumer waits through
PendingDataValueTaskSource<T>:
- The consumer tries the selected partition.
- If it finds no data, it tries every other assigned partition.
- If none has data, it resets the pending-data value-task source and waits.
- A producer appends data to a partition.
- The partition notifies registered consumers through
IBufferPartitionConsumer<TItem>. - The consumer increments its pending-data version and completes the pending value task.
- The consumer resumes the same ring scan; the notification is a wake-up hint, not a priority override.
The pending-data version prevents a lost wake-up when data arrives between the final pull attempt and the transition into the waiting state.
Push-consumer mode is built on pull consumers. The host service discovers push consumers by attribute, creates the corresponding pull consumers, and passes batches to the push-consumer implementation.
Auto-commit push consumers advance progress after a successful pull and before application code
processes the batch. A handler failure therefore does not make that batch eligible for replay.
Manual-commit push consumers receive IBufferConsumerCommitter and decide when to
commit; an uncommitted batch may be delivered again.
The configured ServiceLifetime controls push-consumer resolution:
- A
Singletonconsumer is resolved from the root provider and reused across batches and concurrent consumer loops, so it must be thread-safe. ScopedandTransientconsumers are resolved in a new asynchronous DI scope for every delivered batch. The scope and captured services are asynchronously disposed after the handler completes or throws. Scoped dependencies must not escape that handler call.
Manual commit provides at-least-once delivery:
- A batch that is read but not committed can be delivered again.
- Auto commit advances progress immediately after a successful pull.
- Manual commit advances progress when application code calls
CommitAsync.
Memory keeps offsets only for the lifetime of the process. MemoryMappedFile persists producer offsets and committed consumer offsets to disk. A consumer commit first forces pending MMF log data to a flush boundary. In Batch mode, an uncommitted partial tail batch is not guaranteed to survive an abnormal termination.
See Memory storage and MemoryMappedFile storage for storage-specific behavior.