using RabbitMQ.Client; using RabbitMQ.Client.Events; using Microsoft.Extensions.Configuration; using System.Text; using System.Text.Encodings.Web; using System.Text.Json; namespace QYZH.InteractiveMagazine.Infrastructure.RabbitMQ { public class RabbitMQService : IRabbitMQService { private readonly IRabbitMQConnection _connection; private readonly IConfiguration _configuration; private readonly JsonSerializerOptions options = new JsonSerializerOptions { Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping }; public RabbitMQService(IRabbitMQConnection connection, IConfiguration configuration) { _connection = connection ?? throw new ArgumentNullException(nameof(connection)); _configuration = configuration; } public async Task SendAsync(RabbitMQSendParam param, CancellationToken cancellationToken = default) { try { using var channel = await _connection.CreateChannel(); // 声明 Exchange(持久化) await channel.ExchangeDeclareAsync(exchange: param.Exchange, type: "direct", durable: true, autoDelete: false, arguments: null); // 声明队列(持久化) await channel.QueueDeclareAsync(queue: param.Queue, durable: true, exclusive: false, autoDelete: false, arguments: null); // 绑定队列到 Exchange await channel.QueueBindAsync(queue: param.Queue, exchange: param.Exchange, routingKey: param.RoutingKey, arguments: null); // 清空队列 if (param.Purge) await channel.QueuePurgeAsync(param.Queue); // 消息序列化 var mesjson = JsonSerializer.Serialize(param.Data, options); var body = Encoding.UTF8.GetBytes(mesjson); var properties = new BasicProperties { Persistent = true // 设置消息持久化 }; await channel.BasicPublishAsync(param.Exchange, param.RoutingKey, false, properties, body, cancellationToken); return true; } catch (OperationCanceledException ex) { Console.WriteLine($"Operation was canceled: {ex.Message}"); return false; } catch (Exception ex) { Console.WriteLine($"An error occurred: {ex.Message}"); return false; } } public async Task SendBatchAsync(IEnumerable @params, CancellationToken cancellationToken = default) { IChannel channel = null; try { channel = await _connection.CreateChannel(); // 开启事务 await channel.TxSelectAsync(); var properties = new BasicProperties { Persistent = true // 设置消息持久化 }; // 批量发送消息到不同的 routingKey var declaredExchanges = new HashSet(); foreach (var param in @params) { // 声明 Exchange(持久化) if (declaredExchanges.Add(param.Exchange)) { await channel.ExchangeDeclareAsync(exchange: param.Exchange, type: "direct", durable: true, autoDelete: false, arguments: null); } // 声明队列(持久化) await channel.QueueDeclareAsync(queue: param.Queue, durable: true, exclusive: false, autoDelete: false, arguments: null); // 绑定队列到 Exchange await channel.QueueBindAsync(queue: param.Queue, exchange: param.Exchange, routingKey: param.RoutingKey, arguments: null); // 清空队列 if (param.Purge) await channel.QueuePurgeAsync(param.Queue); // 消息序列化 var mesjson = JsonSerializer.Serialize(param.Data, options); var body = Encoding.UTF8.GetBytes(mesjson); // 发布消息 await channel.BasicPublishAsync(param.Exchange, param.RoutingKey, false, properties, body, cancellationToken); } // 提交事务 - 确保所有消息都发送成功 await channel.TxCommitAsync(); return true; } catch (Exception ex) { Console.WriteLine($"An error occurred: {ex.Message}"); // 回滚事务 try { await channel?.TxRollbackAsync(); } catch { /* 忽略回滚异常 */ } return false; } finally { if (channel != null && !channel.IsClosed) { await channel.CloseAsync(); } } } public async Task ReceiveAsync(string exchange, string queueName, string routingKey, Func callback, CancellationToken cancellationToken = default) { var channel = await _connection.CreateChannel(); var prefetchCount = _configuration.GetValue("RabbitMq:PrefetchCount"); if (prefetchCount == 0) { prefetchCount = 1; } await channel.BasicQosAsync(0, prefetchCount, false, cancellationToken); // 声明 Exchange(持久化) await channel.ExchangeDeclareAsync(exchange: exchange, type: "direct", durable: true, autoDelete: false, arguments: null); // 声明队列(持久化) await channel.QueueDeclareAsync(queue: queueName, durable: true, exclusive: false, autoDelete: false, arguments: null); // 绑定队列到 Exchange await channel.QueueBindAsync(queue: queueName, exchange: exchange, routingKey: routingKey, arguments: null); var consumer = new AsyncEventingBasicConsumer(channel); consumer.ReceivedAsync += async (model, ea) => { //var body = ea.Body.ToArray(); try { // 直接传递 model 和 body 给 callback,不需要转换 await callback(channel, ea); } finally { //await channel.BasicAckAsync(ea.DeliveryTag, false, cancellationToken); } }; await channel.BasicConsumeAsync(queue: queueName, autoAck: false, consumer: consumer, cancellationToken: cancellationToken); // Prevent the method from returning immediately await Task.Delay(-1, cancellationToken); } } }