diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..71c987e --- /dev/null +++ b/.gitignore @@ -0,0 +1,19 @@ +## IDE +.vs/ +.vscode/ +*.suo +*.user +*.userosscache +*.sln.docstates + +## Build +bin/ +obj/ + +## Logs +logs/ +*.log + +## Environment +.env +*.env diff --git a/.trae/documents/.NET_MQ_Consumer_Project_Plan.md b/.trae/documents/.NET_MQ_Consumer_Project_Plan.md new file mode 100644 index 0000000..0996c95 --- /dev/null +++ b/.trae/documents/.NET_MQ_Consumer_Project_Plan.md @@ -0,0 +1,116 @@ +# QuestionLibraryMQConsumer 项目搭建计划 + +## 项目概述 +搭建一个基于 .NET 8 的 RabbitMQ 消费者服务,用于消费题库相关的消息队列消息。 + +## 技术选型 +- **.NET 8** - 最新 LTS 版本 +- **Worker Service** - 后台服务模板,适合 MQ 消费者场景 +- **RabbitMQ.Client 7.x** - 官方 RabbitMQ .NET 客户端 +- **Serilog** - 结构化日志 +- **Microsoft.Extensions.Hosting** - 依赖注入和配置管理 +- **System.Text.Json** - JSON 消息序列化 + +## 项目结构 +``` +QuestionLibraryMQConsumer/ +├── QuestionLibraryMQConsumer.sln +├── src/ +│ └── QuestionLibraryMQConsumer/ +│ ├── QuestionLibraryMQConsumer.csproj +│ ├── Program.cs +│ ├── appsettings.json +│ ├── appsettings.Development.json +│ ├── Consumers/ +│ │ └── QuestionMessageConsumer.cs +│ ├── Services/ +│ │ └── IMessageProcessor.cs +│ │ └── QuestionMessageProcessor.cs +│ ├── Models/ +│ │ └── QuestionMessage.cs +│ ├── Configuration/ +│ │ └── RabbitMqOptions.cs +│ └── Extensions/ +│ └── ServiceCollectionExtensions.cs +└── .gitignore +``` + +## 实施步骤 + +### 第一步: 创建 .NET 8 Worker Service 项目 +1. 创建解决方案文件 +2. 创建 Worker Service 项目 +3. 添加必要的 NuGet 包: + - RabbitMQ.Client + - Serilog.AspNetCore + - Serilog.Sinks.Console + - Serilog.Sinks.File + - Microsoft.Extensions.Hosting + +### 第二步: 创建配置模型 +1. 创建 `RabbitMqOptions.cs` - MQ 连接配置 + - Host: RabbitMQ 服务器地址 + - Port: 端口号 (默认 5672) + - UserName: 用户名 + - Password: 密码 + - VirtualHost: 虚拟主机 + - QueueName: 队列名称 + - ExchangeName: 交换机名称 + - RoutingKey: 路由键 + - PrefetchCount: 预取数量 + - ConsumerCount: 消费者数量 + +### 第三步: 创建消息模型 +1. 创建 `QuestionMessage.cs` - 题库消息数据模型 + - MessageId: 消息 ID + - QuestionId: 题目 ID + - Action: 操作类型 (Add/Update/Delete) + - Content: 消息内容 + - Timestamp: 时间戳 + - 其他题库相关字段 + +### 第三步: 创建消息处理器 +1. 创建 `IMessageProcessor.cs` - 消息处理接口 +2. 创建 `QuestionMessageProcessor.cs` - 具体实现 + - 实现消息处理逻辑 + - 支持重试机制 + - 记录处理日志 + +### 第四步: 创建消费者服务 +1. 创建 `QuestionMessageConsumer.cs` - 继承 BackgroundService + - 实现 RabbitMQ 连接和通道管理 + - 注册消费者 + - 处理消息投递 + - 手动 ACK 确认 + - 处理异常和死信 + - 实现优雅关闭 + +### 第五步: 配置依赖注入 +1. 创建 `ServiceCollectionExtensions.cs` + - 注册 RabbitMQ 配置 + - 注册消息处理器 + - 注册消费者服务 + +### 第六步: 创建主程序入口 +1. 修改 `Program.cs` + - 配置主机 + - 加载配置 + - 注册服务 + - 配置 Serilog 日志 + +### 第七步: 创建配置文件 +1. 创建 `appsettings.json` - 默认配置 +2. 创建 `appsettings.Development.json` - 开发环境配置 +3. 创建 `.gitignore` - Git 忽略文件 + +### 第八步: 验证项目 +1. 执行 `dotnet build` 确保编译通过 +2. 验证项目结构完整性 + +## 关键设计决策 +1. **手动 ACK** - 确保消息不丢失,处理成功后才确认 +2. **Prefetch 限制** - 设置合理的预取数量,避免内存溢出 +3. **连接恢复** - 使用 RabbitMQ.Client 的自动恢复机制 +4. **优雅关闭** - 在应用停止时正确关闭消费者和连接 +5. **结构化日志** - 使用 Serilog 记录详细日志,便于排查问题 +6. **配置外部化** - 所有 MQ 配置通过配置文件管理 diff --git a/QuestionLibraryMQConsumer.slnx b/QuestionLibraryMQConsumer.slnx new file mode 100644 index 0000000..9d122de --- /dev/null +++ b/QuestionLibraryMQConsumer.slnx @@ -0,0 +1,3 @@ + + + diff --git a/QuestionLibraryMQConsumer/Configuration/RabbitMqOptions.cs b/QuestionLibraryMQConsumer/Configuration/RabbitMqOptions.cs new file mode 100644 index 0000000..a55ee8e --- /dev/null +++ b/QuestionLibraryMQConsumer/Configuration/RabbitMqOptions.cs @@ -0,0 +1,27 @@ +namespace QuestionLibraryMQConsumer.Configuration; + +public class RabbitMqOptions +{ + public const string SectionName = "RabbitMq"; + + public string UserName { get; set; } = "guest"; + public string Password { get; set; } = string.Empty; + public string HostName { get; set; } = "localhost"; + public int Port { get; set; } = 5672; + public string ClientProvidedName { get; set; } = string.Empty; + public string VirtualHost { get; set; } = "/"; + public string ExchangeName { get; set; } = string.Empty; + public ushort PrefetchCount { get; set; } = 10; + public List Queues { get; set; } = new(); +} + +public class ConsumerQueueOptions +{ + public string QueueName { get; set; } = string.Empty; + public string RoutingKey { get; set; } = string.Empty; + public string? ExchangeName { get; set; } + public string HandlerType { get; set; } = string.Empty; + public bool Durable { get; set; } = true; + public bool Exclusive { get; set; } = false; + public bool AutoDelete { get; set; } = false; +} diff --git a/QuestionLibraryMQConsumer/Consumers/MultiQueueConsumerService.cs b/QuestionLibraryMQConsumer/Consumers/MultiQueueConsumerService.cs new file mode 100644 index 0000000..d477d42 --- /dev/null +++ b/QuestionLibraryMQConsumer/Consumers/MultiQueueConsumerService.cs @@ -0,0 +1,191 @@ +using System.Text; +using Microsoft.Extensions.Options; +using RabbitMQ.Client; +using RabbitMQ.Client.Events; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; +using QuestionLibraryMQConsumer.Configuration; +using QuestionLibraryMQConsumer.Handlers; + +namespace QuestionLibraryMQConsumer.Consumers; + +public class MultiQueueConsumerService : BackgroundService +{ + private readonly ILogger _logger; + private readonly RabbitMqOptions _rabbitMqOptions; + private readonly IServiceProvider _serviceProvider; + private IConnection? _connection; + private readonly List _channels = new(); + + public MultiQueueConsumerService( + ILogger logger, + IOptions rabbitMqOptions, + IServiceProvider serviceProvider) + { + _logger = logger; + _rabbitMqOptions = rabbitMqOptions.Value; + _serviceProvider = serviceProvider; + } + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + _connection = await CreateConnectionAsync(stoppingToken); + + foreach (var queueConfig in _rabbitMqOptions.Queues) + { + await SetupConsumerAsync(queueConfig, stoppingToken); + } + + stoppingToken.Register(() => + { + _logger.LogInformation("Shutting down RabbitMQ consumer..."); + _connection?.Dispose(); + }); + + await Task.CompletedTask; + } + + private async Task CreateConnectionAsync(CancellationToken stoppingToken) + { + var factory = new ConnectionFactory + { + HostName = _rabbitMqOptions.HostName, + Port = _rabbitMqOptions.Port, + UserName = _rabbitMqOptions.UserName, + Password = _rabbitMqOptions.Password, + VirtualHost = _rabbitMqOptions.VirtualHost + }; + + if (!string.IsNullOrEmpty(_rabbitMqOptions.ClientProvidedName)) + { + factory.ClientProvidedName = _rabbitMqOptions.ClientProvidedName; + } + + var connection = await factory.CreateConnectionAsync(stoppingToken); + _logger.LogInformation("Connected to RabbitMQ at {HostName}:{Port}", + _rabbitMqOptions.HostName, _rabbitMqOptions.Port); + + return connection; + } + + private async Task SetupConsumerAsync(ConsumerQueueOptions queueConfig, CancellationToken stoppingToken) + { + var connection = _connection ?? throw new InvalidOperationException("RabbitMQ connection is not initialized"); + var channel = await connection.CreateChannelAsync(cancellationToken: stoppingToken); + _channels.Add(channel); + + await channel.BasicQosAsync(0, _rabbitMqOptions.PrefetchCount, false, stoppingToken); + + var exchangeName = !string.IsNullOrEmpty(queueConfig.ExchangeName) + ? queueConfig.ExchangeName + : _rabbitMqOptions.ExchangeName; + + await channel.ExchangeDeclareAsync( + exchange: exchangeName, + type: ExchangeType.Direct, + durable: true, + autoDelete: false, + arguments: null, + cancellationToken: stoppingToken); + + await channel.QueueDeclareAsync( + queue: queueConfig.QueueName, + durable: queueConfig.Durable, + exclusive: queueConfig.Exclusive, + autoDelete: queueConfig.AutoDelete, + arguments: null, + cancellationToken: stoppingToken); + + await channel.QueueBindAsync( + queue: queueConfig.QueueName, + exchange: exchangeName, + routingKey: queueConfig.RoutingKey, + arguments: null, + cancellationToken: stoppingToken); + + var consumer = new AsyncEventingBasicConsumer(channel); + consumer.ReceivedAsync += async (sender, eventArgs) => + { + try + { + _logger.LogInformation("[事件触发] ReceivedAsync事件触发,开始处理消息"); + await HandleMessageAsync(channel, eventArgs, queueConfig, stoppingToken); + } + catch (Exception ex) + { + _logger.LogError(ex, "[事件处理] ReceivedAsync事件处理器捕获到未处理异常"); + } + }; + + await channel.BasicConsumeAsync( + queue: queueConfig.QueueName, + autoAck: false, + consumer: consumer, + cancellationToken: stoppingToken); + + _logger.LogInformation("Started consumer for queue {QueueName} on exchange {ExchangeName} with routing key {RoutingKey}", + queueConfig.QueueName, exchangeName, queueConfig.RoutingKey); + } + + private async Task HandleMessageAsync( + IChannel channel, + BasicDeliverEventArgs eventArgs, + ConsumerQueueOptions queueConfig, + CancellationToken stoppingToken) + { + _logger.LogInformation("[步骤1] 开始HandleMessageAsync方法"); + + var body = eventArgs.Body.ToArray(); + var deliveryTag = eventArgs.DeliveryTag; + var messageJson = Encoding.UTF8.GetString(body); + + _logger.LogInformation("[步骤2] 消息解析完成,Queue={QueueName}, DeliveryTag={DeliveryTag}", + queueConfig.QueueName, deliveryTag); + + try + { + _logger.LogInformation("[步骤3] 开始创建DI Scope"); + using var scope = _serviceProvider.CreateScope(); + _logger.LogInformation("[步骤4] DI Scope创建成功"); + + _logger.LogInformation("[步骤5] 开始获取MessageHandlerRegistry"); + var handlerRegistry = scope.ServiceProvider.GetRequiredService(); + _logger.LogInformation("[步骤6] MessageHandlerRegistry获取成功"); + + var registeredKeys = handlerRegistry.GetRegisteredRoutingKeys().ToList(); + _logger.LogInformation("[步骤7] 已注册的RoutingKeys: {Keys}", string.Join(", ", registeredKeys)); + + _logger.LogInformation("[步骤8] 开始查找Handler, RoutingKey={RoutingKey}", queueConfig.RoutingKey); + var handler = handlerRegistry.GetHandler(queueConfig.RoutingKey); + + if (handler == null) + { + _logger.LogWarning("[Handler查找] 未找到Handler for routing key {RoutingKey}", queueConfig.RoutingKey); + await channel.BasicRejectAsync(deliveryTag, false, CancellationToken.None); + return; + } + + _logger.LogInformation("[步骤9] Handler查找成功: {HandlerType}", handler.GetType().Name); + + if (handler.RequiresManualAck) + { + _logger.LogInformation("[步骤10] 调用HandleWithManualAckAsync"); + await handler.HandleWithManualAckAsync(messageJson, channel, deliveryTag, CancellationToken.None); + } + else + { + _logger.LogInformation("[步骤10] 调用HandleJsonAsync"); + await handler.HandleJsonAsync(messageJson, CancellationToken.None); + _logger.LogInformation("[步骤11] HandleJsonAsync执行完成"); + await channel.BasicAckAsync(deliveryTag, false, CancellationToken.None); + _logger.LogInformation("[消息确认] Message processed and acknowledged on queue {QueueName}", queueConfig.QueueName); + } + } + catch (Exception ex) + { + _logger.LogError(ex, "[消息异常] Error processing message on queue {QueueName} with delivery tag {DeliveryTag}", + queueConfig.QueueName, deliveryTag); + await channel.BasicNackAsync(deliveryTag, false, false, CancellationToken.None); + } + } +} diff --git a/QuestionLibraryMQConsumer/Consumers/QuestionMessageConsumer.cs b/QuestionLibraryMQConsumer/Consumers/QuestionMessageConsumer.cs new file mode 100644 index 0000000..a2cbee3 --- /dev/null +++ b/QuestionLibraryMQConsumer/Consumers/QuestionMessageConsumer.cs @@ -0,0 +1,28 @@ +using System.Text; +using System.Text.Json; +using Microsoft.Extensions.Options; +using RabbitMQ.Client; +using RabbitMQ.Client.Events; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; +using QuestionLibraryMQConsumer.Configuration; +using QuestionLibraryMQConsumer.Models; + +namespace QuestionLibraryMQConsumer.Consumers; + +public class QuestionMessageConsumer : BackgroundService +{ + private readonly ILogger _logger; + + public QuestionMessageConsumer( + ILogger logger) + { + _logger = logger; + } + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + _logger.LogInformation("QuestionMessageConsumer is deprecated, use MultiQueueConsumerService instead"); + await Task.CompletedTask; + } +} diff --git a/QuestionLibraryMQConsumer/Data/QuestionLibraryDb.cs b/QuestionLibraryMQConsumer/Data/QuestionLibraryDb.cs new file mode 100644 index 0000000..6eb73c6 --- /dev/null +++ b/QuestionLibraryMQConsumer/Data/QuestionLibraryDb.cs @@ -0,0 +1,26 @@ +using SqlSugar; +using QuestionLibraryMQConsumer.Entities; + +namespace QuestionLibraryMQConsumer.Data; + +public class QuestionLibraryDb +{ + public SqlSugarClient Db { get; private set; } + + public QuestionLibraryDb(string connectionString) + { + Db = new SqlSugarClient(new ConnectionConfig + { + ConnectionString = connectionString, + DbType = DbType.MySql, + IsAutoCloseConnection = true, + InitKeyType = InitKeyType.Attribute, + MoreSettings = new ConnMoreSettings + { + SqlServerCodeFirstNvarchar = true + } + }); + } + + public ISugarQueryable Books => Db.Queryable(); +} diff --git a/QuestionLibraryMQConsumer/Entities/Book.cs b/QuestionLibraryMQConsumer/Entities/Book.cs new file mode 100644 index 0000000..b41c16d --- /dev/null +++ b/QuestionLibraryMQConsumer/Entities/Book.cs @@ -0,0 +1,200 @@ +using SqlSugar; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace QuestionLibraryMQConsumer.Entities +{ + /// + ///书籍 + /// + [SugarTable("book")] + public partial class Book : SqlSugarBaseEntity + { + /// + /// Desc:主键id + /// Default: + /// Nullable:False + /// + [SugarColumn(IsPrimaryKey = true, ColumnName = "id")] + public long Id { get; set; } + + /// + /// Desc:书籍名称 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "name")] + public string Name { get; set; } + + /// + /// Desc:总页数 + /// Default: + /// Nullable:False + /// + [SugarColumn(ColumnName = "totalpage")] + public int TotalPage { get; set; } + + /// + /// Desc:校验页数 + /// Default: + /// Nullable:False + /// + [SugarColumn(ColumnName = "verifypage")] + public int VerifyPage { get; set; } + + /// + /// Desc:状态 + /// Default: + /// Nullable:False + /// + [SugarColumn(ColumnName = "status")] + public int Status { get; set; } + + /// + /// Desc:pdf地址 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "pdfurl")] + public string PdfUrl { get; set; } + + /// + /// Desc:宽度 + /// Default: + /// Nullable:False + /// + [SugarColumn(ColumnName = "width")] + public float Width { get; set; } + + /// + /// Desc:高度 + /// Default: + /// Nullable:False + /// + [SugarColumn(ColumnName = "height")] + public float Height { get; set; } + + /// + /// Desc:乐观锁 + /// Default: + /// Nullable:False + /// + [SugarColumn(ColumnName = "revision")] + public int Revision { get; set; } + + /// + /// Desc:审查结果(默认0 1通过 -1驳回) + /// Default:0 + /// Nullable:False + /// + [SugarColumn(ColumnName = "result")] + public int Result { get; set; } + + /// + /// Desc:封面 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "cover")] + public string Cover { get; set; } + + /// + /// Desc:封底 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "backcover")] + public string BackCover { get; set; } + + /// + /// Desc:pdf预览地址 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "pdfpreviewurl")] + public string PdfPreviewUrl { get; set; } + + /// + /// Desc:副标题 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "title")] + public string Title { get; set; } + + /// + /// Desc:下载书籍页码点阵码PDF文件名称 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "downloadbookpagepdfname")] + public string Downloadbookpagepdfname { get; set; } + + /// + /// Desc:是否删除 + /// Default:b'0' + /// Nullable:False + /// + [SugarColumn(ColumnName = "is_deleted")] + public bool IsDeleted { get; set; } + + /// + /// Desc: + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "created_by")] + public long? CreatedBy { get; set; } + + /// + /// Desc:创建时间 + /// Default: + /// Nullable:False + /// + [SugarColumn(ColumnName = "created_time")] + public DateTime CreatedTime { get; set; } + + /// + /// Desc: + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "updated_by")] + public long? UpdatedBy { get; set; } + + /// + /// Desc:修改时间 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "updated_time")] + public DateTime? UpdatedTime { get; set; } + + /// + /// Desc:书籍描述/简介 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "description")] + public string Description { get; set; } + + /// + /// Desc:书籍类型 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "type")] + public string Type { get; set; } + + /// + /// Desc:书籍来源 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "source")] + public string Source { get; set; } + } +} diff --git a/QuestionLibraryMQConsumer/Entities/SqlSugarBaseEntity.cs b/QuestionLibraryMQConsumer/Entities/SqlSugarBaseEntity.cs new file mode 100644 index 0000000..c1143b8 --- /dev/null +++ b/QuestionLibraryMQConsumer/Entities/SqlSugarBaseEntity.cs @@ -0,0 +1,61 @@ +using SqlSugar; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace QuestionLibraryMQConsumer.Entities +{ + public class SqlSugarBaseEntity + { + /// + /// Desc:是否删除 + /// Default:b'0' + /// Nullable:False + /// + [SugarColumn(ColumnName = "is_deleted")] + public bool IsDeleted { get; set; } = false; + + /// + /// Desc:创建人 + /// Default: + /// Nullable:False + /// + [SugarColumn(ColumnName = "created_by", IsOnlyIgnoreUpdate = true)] + public long CreatedBy { get; set; } + + /// + /// Desc:创建时间 + /// Default: + /// Nullable:False + /// + [SugarColumn(ColumnName = "created_time", IsOnlyIgnoreUpdate = true)] + public DateTime CreatedTime { get; set; } + + /// + /// Desc:修改人 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "updated_by", IsOnlyIgnoreInsert = true)] + public long UpdatedBy { get; set; } + + /// + /// Desc:修改时间 + /// Default: + /// Nullable:True + /// + [SugarColumn(ColumnName = "updated_time", IsOnlyIgnoreInsert = true)] + public DateTime? UpdatedTime { get; set; } + } + + public class SqlSugarBaseEntity : SqlSugarBaseEntity where TKey : struct + { + /// + /// 主键Id + /// + [SugarColumn(IsPrimaryKey = true, ColumnName = "id")] + public TKey Id { get; set; } + } +} diff --git a/QuestionLibraryMQConsumer/Extensions/ServiceCollectionExtensions.cs b/QuestionLibraryMQConsumer/Extensions/ServiceCollectionExtensions.cs new file mode 100644 index 0000000..7f8428f --- /dev/null +++ b/QuestionLibraryMQConsumer/Extensions/ServiceCollectionExtensions.cs @@ -0,0 +1,49 @@ +using Microsoft.Extensions.Configuration; +using Microsoft.Extensions.DependencyInjection; +using QuestionLibraryMQConsumer.Configuration; +using QuestionLibraryMQConsumer.Consumers; +using QuestionLibraryMQConsumer.Data; +using QuestionLibraryMQConsumer.Handlers; + +namespace QuestionLibraryMQConsumer.Extensions; + +public static class ServiceCollectionExtensions +{ + public static IServiceCollection AddRabbitMqConsumer( + this IServiceCollection services, + IConfiguration configuration) + { + services.Configure( + configuration.GetSection(RabbitMqOptions.SectionName)); + + var connectionString = configuration.GetConnectionString("zybDb") ?? throw new InvalidOperationException("Connection string 'zybDb' is not configured"); + services.AddSingleton(new QuestionLibraryDb(connectionString)); + + services.AddSingleton(); + + services.AddTransient(); + + services.AddHostedService(); + + return services; + } + + public static IServiceProvider RegisterMessageHandlers(this IServiceProvider serviceProvider) + { + var registry = serviceProvider.GetRequiredService(); + var handlers = serviceProvider.GetServices().ToList(); + + Console.WriteLine($"[Handler注册] 开始注册Handler,共发现 {handlers.Count} 个Handler"); + + foreach (var handler in handlers) + { + registry.Register(handler); + Console.WriteLine($"[Handler注册] 已注册: {handler.GetType().Name}, RoutingKey: {handler.RoutingKey}"); + } + + var registeredKeys = registry.GetRegisteredRoutingKeys().ToList(); + Console.WriteLine($"[Handler注册] 已注册的RoutingKeys: {string.Join(", ", registeredKeys)}"); + + return serviceProvider; + } +} diff --git a/QuestionLibraryMQConsumer/Handlers/BookCreatedHandler.cs b/QuestionLibraryMQConsumer/Handlers/BookCreatedHandler.cs new file mode 100644 index 0000000..26eb559 --- /dev/null +++ b/QuestionLibraryMQConsumer/Handlers/BookCreatedHandler.cs @@ -0,0 +1,109 @@ +using QuestionLibraryMQConsumer.Data; +using QuestionLibraryMQConsumer.Entities; +using QuestionLibraryMQConsumer.Models; +using RabbitMQ.Client; +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Text.Json; +using System.Threading.Tasks; + +namespace QuestionLibraryMQConsumer.Handlers +{ + public class BookCreatedHandler : IMessageHandler + { + private static readonly JsonSerializerOptions _jsonOptions = new JsonSerializerOptions + { + PropertyNameCaseInsensitive = true + }; + + private readonly ILogger _logger; + private readonly QuestionLibraryDb _db; + + public string RoutingKey => "questionLibrary.book.created"; + public Type MessageType => typeof(BookCreatedMessage); + public bool RequiresManualAck => false; + + public BookCreatedHandler( + ILogger logger, + QuestionLibraryDb db) + { + _logger = logger; + _db = db; + } + + public async Task HandleJsonAsync(string messageJson, CancellationToken cancellationToken) + { + _logger.LogInformation("[Handler入口] HandleJsonAsync被调用, messageJson长度: {Length}", messageJson?.Length ?? 0); + + var message = JsonSerializer.Deserialize(messageJson, _jsonOptions); + + if (message == null) + { + _logger.LogError("[Handler错误] 反序列化失败, messageJson: {Json}", messageJson); + throw new JsonException("Failed to deserialize BookCreatedMessage"); + } + + _logger.LogInformation("[Handler反序列化] 成功, MessageId: {MessageId}, BookId: {BookId}", message.MessageId, message.BookId); + await HandleAsync(message, cancellationToken); + } + + public async Task HandleAsync(BookCreatedMessage message, CancellationToken cancellationToken) + { + _logger.LogInformation("处理书籍创建消息: {MessageId}, 书籍ID: {BookId}, 书籍名称: {BookName}", + message.MessageId, message.BookId, message.BookName); + + var exists = await _db.Db.Queryable() + .AnyAsync(b => b.Id == message.BookId); + + if (exists) + { + _logger.LogWarning("书籍 {BookId} 已存在,跳过创建", message.BookId); + return; + } + + var book = new Book + { + Id = message.BookId, + Name = message.BookName, + Title = message.Subtitle, + Width = message.PaperWidth, + Height = message.PaperHeight, + Cover = message.CoverUrl, + PdfUrl = message.PdfUrl, + BackCover = message.BackCoverUrl, + CreatedTime = DateTime.UtcNow, + UpdatedTime = DateTime.UtcNow, + Source = message.Source, + Status = (int)BookStatusEnum.Created + }; + + await _db.Db.Insertable(book).ExecuteCommandAsync(); + + _logger.LogInformation("书籍 {BookId} 创建成功", message.BookId); + } + + public async Task HandleWithManualAckAsync(string messageJson, IChannel channel, ulong deliveryTag, CancellationToken cancellationToken) + { + var message = JsonSerializer.Deserialize(messageJson, _jsonOptions); + + if (message == null) + throw new JsonException("Failed to deserialize BookCreatedMessage"); + + try + { + await HandleAsync(message, cancellationToken); + + await channel.BasicAckAsync(deliveryTag, false, CancellationToken.None); + _logger.LogInformation("消息 {MessageId} 处理成功并已确认", message.MessageId); + } + catch (Exception ex) + { + _logger.LogError(ex, "消息 {MessageId} 处理失败,拒绝消息", message.MessageId); + await channel.BasicNackAsync(deliveryTag, false, false, CancellationToken.None); + } + } + + } +} diff --git a/QuestionLibraryMQConsumer/Handlers/IMessageHandler.cs b/QuestionLibraryMQConsumer/Handlers/IMessageHandler.cs new file mode 100644 index 0000000..1d1e0c9 --- /dev/null +++ b/QuestionLibraryMQConsumer/Handlers/IMessageHandler.cs @@ -0,0 +1,16 @@ +using RabbitMQ.Client; +using RabbitMQ.Client.Events; + +namespace QuestionLibraryMQConsumer.Handlers; + +public interface IMessageHandler +{ + string RoutingKey { get; } + Type MessageType { get; } + + Task HandleJsonAsync(string messageJson, CancellationToken cancellationToken); + + Task HandleWithManualAckAsync(string messageJson, IChannel channel, ulong deliveryTag, CancellationToken cancellationToken); + + bool RequiresManualAck { get; } +} diff --git a/QuestionLibraryMQConsumer/Handlers/MessageHandlerRegistry.cs b/QuestionLibraryMQConsumer/Handlers/MessageHandlerRegistry.cs new file mode 100644 index 0000000..1f81b64 --- /dev/null +++ b/QuestionLibraryMQConsumer/Handlers/MessageHandlerRegistry.cs @@ -0,0 +1,21 @@ +namespace QuestionLibraryMQConsumer.Handlers; + +public class MessageHandlerRegistry +{ + private readonly Dictionary _handlers = new(); + + public void Register(IMessageHandler handler) + { + _handlers[handler.RoutingKey] = handler; + } + + public IMessageHandler? GetHandler(string routingKey) + { + return _handlers.TryGetValue(routingKey, out var handler) ? handler : null; + } + + public IEnumerable GetRegisteredRoutingKeys() + { + return _handlers.Keys; + } +} diff --git a/QuestionLibraryMQConsumer/Models/BookCreatedMessage.cs b/QuestionLibraryMQConsumer/Models/BookCreatedMessage.cs new file mode 100644 index 0000000..0af7106 --- /dev/null +++ b/QuestionLibraryMQConsumer/Models/BookCreatedMessage.cs @@ -0,0 +1,45 @@ +namespace QuestionLibraryMQConsumer.Models; + +public class BookCreatedMessage +{ + /// + /// ��ϢId + /// + public string MessageId { get; set; } = string.Empty; + /// + /// BookId + /// + public long BookId { get; set; } + /// + /// BookName + /// + public string BookName { get; set; } = string.Empty; + /// + /// ������ + /// + public string Subtitle { get; set; } = string.Empty; + /// + /// ֽ�ſ��� + /// + public float PaperWidth { get; set; } + /// + /// ֽ�Ÿ߶� + /// + public float PaperHeight { get; set; } + /// + /// ���� + /// + public string? CoverUrl { get; set; } + /// + /// ��� + /// + public string? BackCoverUrl { get; set; } + /// + /// �鼮PDF + /// + public string? PdfUrl { get; set; } + /// + /// ԴʶжϻشMQϢַ + /// + public string? Source { get; set; } +} diff --git a/QuestionLibraryMQConsumer/Models/BookStatusEnum.cs b/QuestionLibraryMQConsumer/Models/BookStatusEnum.cs new file mode 100644 index 0000000..95969c4 --- /dev/null +++ b/QuestionLibraryMQConsumer/Models/BookStatusEnum.cs @@ -0,0 +1,60 @@ +using System; +using System.Collections.Generic; +using System.ComponentModel; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace QuestionLibraryMQConsumer.Models +{ + public enum BookStatusEnum + { + /// + /// 已创建 + /// + [Description("已创建")] + Created = 0, + + /// + /// 编辑 + /// + [Description("编辑")] + Editor = 1, + + /// + /// 已校验 + /// + [Description("已校验")] + Verify = 3, + + /// + /// 铺码中 + /// + [Description("铺码中")] + Codeing = 4, + + /// + /// 铺码成功 + /// + [Description("铺码成功")] + CodeSuccess = 5, + + /// + /// 铺码失败 + /// + [Description("铺码失败")] + CodeFail = -5, + + /// + /// 已归档 + /// + [Description("已归档")] + Archive = 9, + + /// + /// 已废弃 + /// + [Description("已废弃")] + Abandoned = -9, + } +} diff --git a/QuestionLibraryMQConsumer/Program.cs b/QuestionLibraryMQConsumer/Program.cs new file mode 100644 index 0000000..da6381e --- /dev/null +++ b/QuestionLibraryMQConsumer/Program.cs @@ -0,0 +1,43 @@ +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using Serilog; +using QuestionLibraryMQConsumer.Extensions; + +namespace QuestionLibraryMQConsumer; + +public class Program +{ + public static void Main(string[] args) + { + Log.Logger = new LoggerConfiguration() + .MinimumLevel.Information() + .WriteTo.Console() + .WriteTo.File("logs/questionlibrary-.txt", rollingInterval: RollingInterval.Day) + .CreateLogger(); + + try + { + Log.Information("Starting QuestionLibrary MQ Consumer"); + + var host = CreateHostBuilder(args).Build(); + host.Services.RegisterMessageHandlers(); + host.Run(); + } + catch (Exception ex) + { + Log.Fatal(ex, "Host terminated unexpectedly"); + } + finally + { + Log.CloseAndFlush(); + } + } + + public static IHostBuilder CreateHostBuilder(string[] args) => + Host.CreateDefaultBuilder(args) + .UseSerilog() + .ConfigureServices((hostContext, services) => + { + services.AddRabbitMqConsumer(hostContext.Configuration); + }); +} diff --git a/QuestionLibraryMQConsumer/Properties/launchSettings.json b/QuestionLibraryMQConsumer/Properties/launchSettings.json new file mode 100644 index 0000000..d16f1c5 --- /dev/null +++ b/QuestionLibraryMQConsumer/Properties/launchSettings.json @@ -0,0 +1,12 @@ +{ + "$schema": "http://json.schemastore.org/launchsettings.json", + "profiles": { + "QuestionLibraryMQConsumer": { + "commandName": "Project", + "dotnetRunMessages": true, + "environmentVariables": { + "DOTNET_ENVIRONMENT": "Development" + } + } + } +} diff --git a/QuestionLibraryMQConsumer/QuestionLibraryMQConsumer.csproj b/QuestionLibraryMQConsumer/QuestionLibraryMQConsumer.csproj new file mode 100644 index 0000000..e814550 --- /dev/null +++ b/QuestionLibraryMQConsumer/QuestionLibraryMQConsumer.csproj @@ -0,0 +1,16 @@ + + + + net8.0 + enable + enable + dotnet-QuestionLibraryMQConsumer-3847697e-8fac-49a3-9723-51eb362c06ff + + + + + + + + + diff --git a/QuestionLibraryMQConsumer/appsettings.Development.json b/QuestionLibraryMQConsumer/appsettings.Development.json new file mode 100644 index 0000000..dbd5736 --- /dev/null +++ b/QuestionLibraryMQConsumer/appsettings.Development.json @@ -0,0 +1,56 @@ +{ + "RabbitMq": { + "UserName": "smartschool", + "Password": "@ss%&*otz%d*pq2S", + "HostName": "192.168.20.150", + "Port": 5672, + "ClientProvidedName": "QuestionLibraryMQConsumer", + "VirtualHost": "local", + "ExchangeName": "question.library.exchange", + "PrefetchCount": 10, + "Queues": [ + { + "QueueName": "questionLibrary.book.created.queue", + "ExchangeName": "icr.direct", + "RoutingKey": "questionLibrary.book.created", + "Durable": true, + "Exclusive": false, + "AutoDelete": false + }, + { + "QueueName": "questionLibrary.book.created.queue", + "ExchangeName": "question.library.exchange.dev", + "RoutingKey": "questionLibrary.book.created", + "Durable": true, + "Exclusive": false, + "AutoDelete": false + }, + { + "QueueName": "questionLibrary.book.created.queue", + "ExchangeName": "question.library.exchange", + "RoutingKey": "questionLibrary.book.created", + "Durable": true, + "Exclusive": false, + "AutoDelete": false + } + ] + }, + "ConnectionStrings": { + "zybDb": "server=192.168.20.150;port=13306;database=questionlibrary;user=user;password=n68792bu!y99r905;charset=utf8mb4;" + }, + "Serilog": { + "MinimumLevel": "Information", + "WriteTo": [ + { + "Name": "Console" + }, + { + "Name": "File", + "Args": { + "path": "logs/questionlibrary-.txt", + "rollingInterval": "Day" + } + } + ] + } +} diff --git a/QuestionLibraryMQConsumer/appsettings.json b/QuestionLibraryMQConsumer/appsettings.json new file mode 100644 index 0000000..dbd5736 --- /dev/null +++ b/QuestionLibraryMQConsumer/appsettings.json @@ -0,0 +1,56 @@ +{ + "RabbitMq": { + "UserName": "smartschool", + "Password": "@ss%&*otz%d*pq2S", + "HostName": "192.168.20.150", + "Port": 5672, + "ClientProvidedName": "QuestionLibraryMQConsumer", + "VirtualHost": "local", + "ExchangeName": "question.library.exchange", + "PrefetchCount": 10, + "Queues": [ + { + "QueueName": "questionLibrary.book.created.queue", + "ExchangeName": "icr.direct", + "RoutingKey": "questionLibrary.book.created", + "Durable": true, + "Exclusive": false, + "AutoDelete": false + }, + { + "QueueName": "questionLibrary.book.created.queue", + "ExchangeName": "question.library.exchange.dev", + "RoutingKey": "questionLibrary.book.created", + "Durable": true, + "Exclusive": false, + "AutoDelete": false + }, + { + "QueueName": "questionLibrary.book.created.queue", + "ExchangeName": "question.library.exchange", + "RoutingKey": "questionLibrary.book.created", + "Durable": true, + "Exclusive": false, + "AutoDelete": false + } + ] + }, + "ConnectionStrings": { + "zybDb": "server=192.168.20.150;port=13306;database=questionlibrary;user=user;password=n68792bu!y99r905;charset=utf8mb4;" + }, + "Serilog": { + "MinimumLevel": "Information", + "WriteTo": [ + { + "Name": "Console" + }, + { + "Name": "File", + "Args": { + "path": "logs/questionlibrary-.txt", + "rollingInterval": "Day" + } + } + ] + } +}