Вот сервис: < /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
}
}
}
}
Код: Выделить всё
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);
});
}
Подробнее здесь: https://stackoverflow.com/questions/796 ... pplication