/
vshmidt
/
masstransit
Обзор
Документация
Войти
/
vshmidt
/
masstransit
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
develop
src/Transports/MassTransit.EventHubIntegration/EventHubIntegration/PartitionCheckpointData.cs
41 строка
1 KB
denys.kozhevnikov
Consume messages in order for Riders
24 ноя 2023, 16:59
24 ноя 2023, 16:59
ad274c9
Код
Авторство
О чём код?
namespace MassTransit.EventHubIntegration { using System.Threading; using System.Threading.Tasks; using Azure.Messaging.EventHubs.Processor; using Checkpoints; public class PartitionCheckpointData { readonly CancellationTokenSource _cancellationTokenSource; readonly ICheckpointer _checkpointer; readonly PendingConfirmationCollection _pending; public PartitionCheckpointData(ReceiveSettings settings, PendingConfirmationCollection pending) { _cancellationTokenSource = new CancellationTokenSource(); _checkpointer = new BatchCheckpointer(settings, _cancellationTokenSource.Token); _pending = pending; } public Task Pending(ProcessEventArgs eventArgs) { var pendingConfirmation = _pending.Add(eventArgs); return _checkpointer.Pending(pendingConfirmation); } public async Task Close(PartitionClosingEventArgs args) { if (args.Reason != ProcessingStoppedReason.Shutdown) _cancellationTokenSource.Cancel(); await _checkpointer.DisposeAsync().ConfigureAwait(false); LogContext.Info?.Log("Partition: {PartitionId} was closed, reason: {Reason}", args.PartitionId, args.Reason); _cancellationTokenSource.Cancel(); _cancellationTokenSource.Dispose(); } } }