- Загрузка файлов:
Пользователи может загружать файлы, и при успешной загрузке сервер генерирует URL-адреса для этих файлов. Мне нужно опубликовать эти URL-адреса в очереди RabbitMQ. - Создание шаблонов:
Пользователи могут создавать шаблоны, которые также необходимо опубликовано в другой очереди RabbitMQ.
Я хочу объединить сообщения из двух очередей ( URL-адреса файлов и сведения о шаблоне) перед их совместной обработкой. Конечная цель — отправлять электронные письма, содержащие загруженные файлы в виде вложений, при каждом создании шаблона.
Вопросы:
Каков наилучший подход к объединению сообщений из двух очередей RabbitMQ? надежным способом?
Должен ли я использовать выделенную службу обработки, которая обрабатывает сообщения из обеих очередей, или есть лучшая стратегия?
Как я могу гарантировать, что электронные письма будут отправлены только после того, как оба сообщения будут успешно отправлены объединены?
Будем благодарны за любые советы или примеры!
using Microsoft.Extensions.Options;
using Newtonsoft.Json;
using RabbitMQ.Client.Events;
using RabbitMQ.Client;
using System.Text;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
public class AlertMessageConsumerService : BackgroundService
{
private readonly IServiceProvider _serviceProvider;
private readonly RabbitMQSetting _rabbitMqSetting;
private readonly ILogger _logger;
private IConnection _connection;
private IModel _channel;
public AlertMessageConsumerService(IOptions rabbitMqSetting, IServiceProvider serviceProvider, ILogger logger)
{
_rabbitMqSetting = rabbitMqSetting.Value;
_serviceProvider = serviceProvider;
_logger = logger;
var factory = new ConnectionFactory
{
HostName = _rabbitMqSetting.HostName,
UserName = _rabbitMqSetting.UserName,
Password = _rabbitMqSetting.Password
};
_connection = factory.CreateConnection();
_channel = _connection.CreateModel();
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
StartConsuming("alertQueue", stoppingToken);
StartConsuming("anotherQueue", stoppingToken);
await Task.CompletedTask;
}
private void StartConsuming(string queueName, CancellationToken cancellationToken)
{
_channel.QueueDeclare(queue: queueName, durable: false, exclusive: false, autoDelete: false, arguments: null);
var consumer = new EventingBasicConsumer(_channel);
consumer.Received += async (model, ea) =>
{
var body = ea.Body.ToArray();
var message = Encoding.UTF8.GetString(body);
bool processedSuccessfully = false;
try
{
processedSuccessfully = await ProcessMessageAsync(message);
}
catch (Exception ex)
{
_logger.LogError($"Exception occurred while processing message from queue {queueName}: {ex}");
}
if (processedSuccessfully)
{
_channel.BasicAck(deliveryTag: ea.DeliveryTag, multiple: false);
}
else
{
_channel.BasicReject(deliveryTag: ea.DeliveryTag, requeue: true);
}
};
_channel.BasicConsume(queue: queueName, autoAck: false, consumer: consumer);
}
private async Task ProcessMessageAsync(string message)
{
try
{
using (var scope = _serviceProvider.CreateScope())
{
var emailService = scope.ServiceProvider.GetRequiredService();
var vsoService = scope.ServiceProvider.GetRequiredService();
var alertTypeService = scope.ServiceProvider.GetRequiredService();
var s3Service = scope.ServiceProvider.GetRequiredService();
var alertMessage = JsonConvert.DeserializeObject(message);
var uploadMessage = JsonConvert.DeserializeObject(message);
if (alertMessage != null && uploadMessage != null)
{
_logger.LogInformation($"AlertMessage: {JsonConvert.SerializeObject(alertMessage)}");
_logger.LogInformation($"FileUploadMessage: {JsonConvert.SerializeObject(uploadMessage)}");
return true;
}
else
{
_logger.LogWarning("Both AlertMessage and FileUploadMessage must be present in the message.");
return false;
}
}
}
catch (JsonException jsonEx)
{
_logger.LogError($"JSON error processing message: {jsonEx.Message}");
return false;
}
catch (Exception ex)
{
_logger.LogError($"Error processing message: {ex.Message}");
return false;
}
}
public override void Dispose()
{
_channel.Close();
_connection.Close();
base.Dispose();
}
}
Подробнее здесь: https://stackoverflow.com/questions/791 ... with-rabbi