Как эффективно использовать канал RabbitMQ в приложении ASP.NET? ⇐ C#

Место общения программистов C#
Anonymous
Как эффективно использовать канал RabbitMQ в приложении ASP.NET?

Сообщение Anonymous »

У меня есть веб -приложение, где результат вызова конечной точки опубликует сообщение Rabbitmq. Это говорит о том, что объект канала не защищен потоком, но канал должен быть долго жить. Я придумал эту службу Singleton, чтобы создать одно соединение и канал. Я думать это безопасно поток, потому что все в executeAsync ожидает , и, следовательно, никакие два потока не могли бы вызвать await _channel.basicpublishasync одновременно. Является ли мое понимание правильно? < /P>
Вот сервис: < /p>

Код: Выделить всё

public class RabbitMqService : BackgroundService
{
private readonly  int MAX_TRIES = 4;
private readonly List DELAYS = new List { 500, 2000, 3000 }; // initial delay in milliseconds
private IConnection _connection;
private IChannel _channel;
private readonly ILogger _logger;
private readonly IOptions _options;
private readonly IMessageContainerService _messageContainerService;

public RabbitMqService(IOptions options, ILoggerFactory logFactory, IMessageContainerService messageContainerService)
{
_options = options;
_logger = logFactory.CreateLogger();
_messageContainerService = messageContainerService;

_logger.LogInformation("Instantiating RabbitMqService.");
}

private async Task Connect()
{
var integration = _options.Value;
var connectionFactory = new ConnectionFactory
{
HostName = integration.HostName,
Port = integration.Port,
UserName = integration.UserName,
Password = integration.Password,
AutomaticRecoveryEnabled = true //will attempt to auto recconect https://www.rabbitmq.com/client-libraries/dotnet-api-guide#recovery
};

if (null == _connection || !_connection.IsOpen)
{
_logger.LogInformation($"Creating rabbit connection...  {integration.HostName} | {integration.Port} | {integration.UserName}");
try {
_connection = await connectionFactory.CreateConnectionAsync();
} catch (RabbitMQ.Client.Exceptions.BrokerUnreachableException e) {
_logger.LogError($"Exception Creating rabbit connection {e.Message}", e);
await Task.Delay(5000);
await Connect(); // if connection is not opened auth retries will not take place https://www.rabbitmq.com/client-libraries/dotnet-api-guide#recovery-triggers
}
}

if (null == _channel || _channel.IsClosed)
{
_logger.LogInformation($"Creating rabbit channel...");
_channel = await _connection.CreateChannelAsync();
}
}

private async Task Publish(T model, string exchange, string routingKey)
{
if (null == _channel || _channel.IsClosed)
{
await Connect();
}

var activityName = $"{routingKey} send";
using var activity = new Activity(activityName);
activity.SetIdFormat(ActivityIdFormat.W3C);

if (Activity.Current != null)
{
activity.SetParentId(Activity.Current.Id);
foreach (var baggage in Activity.Current.Baggage)
{
activity.AddBaggage(baggage.Key, baggage.Value);
}
}

activity.Start();

var properties = new BasicProperties();

var json = JsonSerializer.Serialize(model);
var body = Encoding.UTF8.GetBytes(json);

await _channel.BasicPublishAsync(
exchange,
routingKey,
true,
properties,
body
);

_logger.LogInformation($"{JsonSerializer.Serialize(model)} dumped to RabbitMQ successfully.");
}

protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
_logger.LogInformation("Starting RabbitMqService ExecuteAsync");

TimeSpan pollInterval = TimeSpan.FromMilliseconds(100);
using PeriodicTimer timer = new PeriodicTimer(pollInterval);

while (!stoppingToken.IsCancellationRequested &&
await timer.WaitForNextTickAsync(stoppingToken))
{
try
{
PublishDto message = _messageContainerService.Dequeue();

while (null != message)
{
await PublishWithRetries(message.model, message.exchange, message.routingKey);
message = _messageContainerService.Dequeue();
//todo remove once we have an indication for how fast these process
if (_messageContainerService.QueueLength() > 5)
{
_logger.LogInformation($"_messageContainerService queue length is {_messageContainerService.QueueLength()}");
}
}
}
catch (Exception ex)
{
_logger.LogError($"Exception in RabbitMqService ExecuteAsync while loop {ex.Message}", ex);
}
}

if (stoppingToken.IsCancellationRequested)
{
_logger.LogInformation("Stopping RabbitMqService Cancellation Requested");

if (null != _channel && _channel.IsOpen)
{
await _channel.CloseAsync();
}

if (null != _connection && _connection.IsOpen)
{
await _connection.CloseAsync();
}

if (null != _channel)
{
await _channel.DisposeAsync();
}

if (null != _connection)
{
await _connection.DisposeAsync();
}
}
}

private async Task PublishWithRetries(object messageModel, string messageExchange, string messageRoutingKey)
{
int attempt = 0;

while (attempt <  MAX_TRIES)
{
try
{
attempt++;
await Publish(messageModel, messageExchange, messageRoutingKey);
}
catch (Exception ex)
{
_logger.LogError($"Publish failed exception {ex.Message}", ex);

if (attempt == MAX_TRIES)
{
_logger.LogError($"Publish failed for {JsonSerializer.Serialize(messageModel)} with exception {ex.Message}", ex);
return;
}

int backoff = DELAYS[attempt - 1];
_logger.LogError($"Waiting {backoff}ms before retrying...");
await Task.Delay(backoff); // Exponential backoff
}
}
}
}
Сообщения поступают из concurrentqueue Это MessageContainerService :

Код: Выделить всё

public class MessageContainerService : IMessageContainerService
{
private readonly ConcurrentQueue
 _messageQueue = new ConcurrentQueue();

public void Enqueue(Object model, string commandsExchange, string routingKey)
{
_messageQueue.Enqueue(new PublishDto()
{
model = model,
exchange = commandsExchange,
routingKey = routingKey
});
}

public PublishDto? Dequeue()
{
var hasMessage = _messageQueue.TryDequeue(out var message);
return hasMessage ? message : null;
}

public long QueueLength()
{
return _messageQueue.Count;
}
}
< /code>
Так добавляются услуги в виде синглтонов: < /p>
builder.Services.AddSingleton();

for (var i = 0; i < builder.Configuration.GetValue("NumberOfRabbitMqServices"); i++)
{
// this will only add one instance so we have to use AddSingleton
// builder.Services.AddHostedService();
builder.Services.AddSingleton(provider =>
{
var logger = provider.GetRequiredService();
var options = provider.GetRequiredService();
var messageContainer = provider.GetRequiredService();
return new RabbitMqService(options, logger, messageContainer);
});
}
Я пытаюсь сделать долгосрочную безопасную службу потока, используя объект канала Rabbitmq.


Подробнее здесь: https://stackoverflow.com/questions/796 ... pplication

Вернуться в «C#»