/
dsrelax
/
Tarantool.Queue.NetCore
Обзор
Документация
Войти
/
dsrelax
/
Tarantool.Queue.NetCore
Код
Запросы
0
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
samples/StandardTarantoolQueue/StandardTarantoolQueueReader/ConsumerService.cs
84 строки
3 KB
Дмитрий
Change samples and description
01 фев 2024, 19:59
01 фев 2024, 19:59
e1f9c7b
Код
Авторство
О чём код?
using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using System.Diagnostics; using Tarantool.Queues; using Tarantool.Queues.Helpers; using Tarantool.Queues.Options; namespace TarantoolReader { public sealed class ConsumerService : IHostedService, IDisposable { private readonly ILogger<ConsumerService> _logger; private readonly ITubeConsumerBuilder _tubeConsumerBuilder; private readonly CancellationTokenSource _cancellationTokenSource = new(); private Task? _mainTask = null; public ConsumerService(ITubeConsumerBuilder tubeProducerBuilder, ILogger<ConsumerService> logger) { _tubeConsumerBuilder = tubeProducerBuilder; _logger = logger; } public void Dispose() { _cancellationTokenSource.Dispose(); _mainTask?.Dispose(); } public Task StartAsync(CancellationToken cancellationToken) { _mainTask = Task.Factory.StartNew(async() => { //var consumer1 = await _tubeConsumerBuilder.Build<UTubeTtlTubeOptions>("queue_test_utubettl", true); var consumer = await _tubeConsumerBuilder.Build<FiFoTubeOptions>("queue_test_fifo"); int taskCount = 0; TimeSpan allTime = TimeSpan.Zero; Stopwatch sw = new(); Console.Clear(); var cursorPosition = Console.GetCursorPosition(); while (!_cancellationTokenSource.IsCancellationRequested) { try { var consumeResult = await consumer.Consume(null, _cancellationTokenSource.Token); if (consumeResult.Task != null) { taskCount++; var writeDate = DateTime.Parse(consumeResult.Task!.TaskData).ToUniversalTime(); allTime += consumeResult.ConsumeDate - writeDate; sw.Start(); await consumeResult.Commit(); sw.Stop(); Console.SetCursorPosition(cursorPosition.Left, cursorPosition.Top); Console.Write($"Read - {taskCount} average message queuing time - {allTime / taskCount}, average task completion time {sw.Elapsed / taskCount}"); } } catch (TaskCanceledException) { } catch (Exception ex) { Console.WriteLine($"Error reading task from queue\n{ex}"); } } }, _cancellationTokenSource.Token); return Task.CompletedTask; } public async Task StopAsync(CancellationToken cancellationToken) { if (!_cancellationTokenSource.IsCancellationRequested) _cancellationTokenSource.Cancel(); if(_mainTask != null) await _mainTask; } } }