添加项目文件。

This commit is contained in:
glz
2026-05-18 17:08:35 +08:00
parent 2ef331627b
commit aa223688a0
20 changed files with 1154 additions and 0 deletions

19
.gitignore vendored Normal file
View File

@ -0,0 +1,19 @@
## IDE
.vs/
.vscode/
*.suo
*.user
*.userosscache
*.sln.docstates
## Build
bin/
obj/
## Logs
logs/
*.log
## Environment
.env
*.env

View File

@ -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 配置通过配置文件管理

View File

@ -0,0 +1,3 @@
<Solution>
<Project Path="QuestionLibraryMQConsumer/QuestionLibraryMQConsumer.csproj" />
</Solution>

View File

@ -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<ConsumerQueueOptions> 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;
}

View File

@ -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<MultiQueueConsumerService> _logger;
private readonly RabbitMqOptions _rabbitMqOptions;
private readonly IServiceProvider _serviceProvider;
private IConnection? _connection;
private readonly List<IChannel> _channels = new();
public MultiQueueConsumerService(
ILogger<MultiQueueConsumerService> logger,
IOptions<RabbitMqOptions> 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<IConnection> 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<MessageHandlerRegistry>();
_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);
}
}
}

View File

@ -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<QuestionMessageConsumer> _logger;
public QuestionMessageConsumer(
ILogger<QuestionMessageConsumer> logger)
{
_logger = logger;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
_logger.LogInformation("QuestionMessageConsumer is deprecated, use MultiQueueConsumerService instead");
await Task.CompletedTask;
}
}

View File

@ -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<Book> Books => Db.Queryable<Book>();
}

View File

@ -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
{
///<summary>
///书籍
///</summary>
[SugarTable("book")]
public partial class Book : SqlSugarBaseEntity
{
/// <summary>
/// Desc:主键id
/// Default:
/// Nullable:False
/// </summary>
[SugarColumn(IsPrimaryKey = true, ColumnName = "id")]
public long Id { get; set; }
/// <summary>
/// Desc:书籍名称
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "name")]
public string Name { get; set; }
/// <summary>
/// Desc:总页数
/// Default:
/// Nullable:False
/// </summary>
[SugarColumn(ColumnName = "totalpage")]
public int TotalPage { get; set; }
/// <summary>
/// Desc:校验页数
/// Default:
/// Nullable:False
/// </summary>
[SugarColumn(ColumnName = "verifypage")]
public int VerifyPage { get; set; }
/// <summary>
/// Desc:状态
/// Default:
/// Nullable:False
/// </summary>
[SugarColumn(ColumnName = "status")]
public int Status { get; set; }
/// <summary>
/// Desc:pdf地址
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "pdfurl")]
public string PdfUrl { get; set; }
/// <summary>
/// Desc:宽度
/// Default:
/// Nullable:False
/// </summary>
[SugarColumn(ColumnName = "width")]
public float Width { get; set; }
/// <summary>
/// Desc:高度
/// Default:
/// Nullable:False
/// </summary>
[SugarColumn(ColumnName = "height")]
public float Height { get; set; }
/// <summary>
/// Desc:乐观锁
/// Default:
/// Nullable:False
/// </summary>
[SugarColumn(ColumnName = "revision")]
public int Revision { get; set; }
/// <summary>
/// Desc:审查结果(默认0 1通过 -1驳回)
/// Default:0
/// Nullable:False
/// </summary>
[SugarColumn(ColumnName = "result")]
public int Result { get; set; }
/// <summary>
/// Desc:封面
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "cover")]
public string Cover { get; set; }
/// <summary>
/// Desc:封底
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "backcover")]
public string BackCover { get; set; }
/// <summary>
/// Desc:pdf预览地址
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "pdfpreviewurl")]
public string PdfPreviewUrl { get; set; }
/// <summary>
/// Desc:副标题
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "title")]
public string Title { get; set; }
/// <summary>
/// Desc:下载书籍页码点阵码PDF文件名称
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "downloadbookpagepdfname")]
public string Downloadbookpagepdfname { get; set; }
/// <summary>
/// Desc:是否删除
/// Default:b'0'
/// Nullable:False
/// </summary>
[SugarColumn(ColumnName = "is_deleted")]
public bool IsDeleted { get; set; }
/// <summary>
/// Desc:
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "created_by")]
public long? CreatedBy { get; set; }
/// <summary>
/// Desc:创建时间
/// Default:
/// Nullable:False
/// </summary>
[SugarColumn(ColumnName = "created_time")]
public DateTime CreatedTime { get; set; }
/// <summary>
/// Desc:
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "updated_by")]
public long? UpdatedBy { get; set; }
/// <summary>
/// Desc:修改时间
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "updated_time")]
public DateTime? UpdatedTime { get; set; }
/// <summary>
/// Desc:书籍描述/简介
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "description")]
public string Description { get; set; }
/// <summary>
/// Desc:书籍类型
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "type")]
public string Type { get; set; }
/// <summary>
/// Desc:书籍来源
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "source")]
public string Source { get; set; }
}
}

View File

@ -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
{
/// <summary>
/// Desc:是否删除
/// Default:b'0'
/// Nullable:False
/// </summary>
[SugarColumn(ColumnName = "is_deleted")]
public bool IsDeleted { get; set; } = false;
/// <summary>
/// Desc:创建人
/// Default:
/// Nullable:False
/// </summary>
[SugarColumn(ColumnName = "created_by", IsOnlyIgnoreUpdate = true)]
public long CreatedBy { get; set; }
/// <summary>
/// Desc:创建时间
/// Default:
/// Nullable:False
/// </summary>
[SugarColumn(ColumnName = "created_time", IsOnlyIgnoreUpdate = true)]
public DateTime CreatedTime { get; set; }
/// <summary>
/// Desc:修改人
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "updated_by", IsOnlyIgnoreInsert = true)]
public long UpdatedBy { get; set; }
/// <summary>
/// Desc:修改时间
/// Default:
/// Nullable:True
/// </summary>
[SugarColumn(ColumnName = "updated_time", IsOnlyIgnoreInsert = true)]
public DateTime? UpdatedTime { get; set; }
}
public class SqlSugarBaseEntity<TKey> : SqlSugarBaseEntity where TKey : struct
{
/// <summary>
/// 主键Id
/// </summary>
[SugarColumn(IsPrimaryKey = true, ColumnName = "id")]
public TKey Id { get; set; }
}
}

View File

@ -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<RabbitMqOptions>(
configuration.GetSection(RabbitMqOptions.SectionName));
var connectionString = configuration.GetConnectionString("zybDb") ?? throw new InvalidOperationException("Connection string 'zybDb' is not configured");
services.AddSingleton<QuestionLibraryDb>(new QuestionLibraryDb(connectionString));
services.AddSingleton<MessageHandlerRegistry>();
services.AddTransient<IMessageHandler, BookCreatedHandler>();
services.AddHostedService<MultiQueueConsumerService>();
return services;
}
public static IServiceProvider RegisterMessageHandlers(this IServiceProvider serviceProvider)
{
var registry = serviceProvider.GetRequiredService<MessageHandlerRegistry>();
var handlers = serviceProvider.GetServices<IMessageHandler>().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;
}
}

View File

@ -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<BookCreatedHandler> _logger;
private readonly QuestionLibraryDb _db;
public string RoutingKey => "questionLibrary.book.created";
public Type MessageType => typeof(BookCreatedMessage);
public bool RequiresManualAck => false;
public BookCreatedHandler(
ILogger<BookCreatedHandler> 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<BookCreatedMessage>(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<Book>()
.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<BookCreatedMessage>(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);
}
}
}
}

View File

@ -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; }
}

View File

@ -0,0 +1,21 @@
namespace QuestionLibraryMQConsumer.Handlers;
public class MessageHandlerRegistry
{
private readonly Dictionary<string, IMessageHandler> _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<string> GetRegisteredRoutingKeys()
{
return _handlers.Keys;
}
}

View File

@ -0,0 +1,45 @@
namespace QuestionLibraryMQConsumer.Models;
public class BookCreatedMessage
{
/// <summary>
/// <20><>ϢId
/// </summary>
public string MessageId { get; set; } = string.Empty;
/// <summary>
/// BookId
/// </summary>
public long BookId { get; set; }
/// <summary>
/// BookName
/// </summary>
public string BookName { get; set; } = string.Empty;
/// <summary>
/// <20><><EFBFBD><EFBFBD><EFBFBD><EFBFBD>
/// </summary>
public string Subtitle { get; set; } = string.Empty;
/// <summary>
/// ֽ<>ſ<EFBFBD><C5BF><EFBFBD>
/// </summary>
public float PaperWidth { get; set; }
/// <summary>
/// ֽ<>Ÿ߶<C5B8>
/// </summary>
public float PaperHeight { get; set; }
/// <summary>
/// <20><><EFBFBD><EFBFBD>
/// </summary>
public string? CoverUrl { get; set; }
/// <summary>
/// <20><><EFBFBD>
/// </summary>
public string? BackCoverUrl { get; set; }
/// <summary>
/// <20>鼮PDF
/// </summary>
public string? PdfUrl { get; set; }
/// <summary>
/// ԴʶжϻشMQϢַ
/// </summary>
public string? Source { get; set; }
}

View File

@ -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
{
/// <summary>
/// 已创建
/// </summary>
[Description("已创建")]
Created = 0,
/// <summary>
/// 编辑
/// </summary>
[Description("编辑")]
Editor = 1,
/// <summary>
/// 已校验
/// </summary>
[Description("已校验")]
Verify = 3,
/// <summary>
/// 铺码中
/// </summary>
[Description("铺码中")]
Codeing = 4,
/// <summary>
/// 铺码成功
/// </summary>
[Description("铺码成功")]
CodeSuccess = 5,
/// <summary>
/// 铺码失败
/// </summary>
[Description("铺码失败")]
CodeFail = -5,
/// <summary>
/// 已归档
/// </summary>
[Description("已归档")]
Archive = 9,
/// <summary>
/// 已废弃
/// </summary>
[Description("已废弃")]
Abandoned = -9,
}
}

View File

@ -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);
});
}

View File

@ -0,0 +1,12 @@
{
"$schema": "http://json.schemastore.org/launchsettings.json",
"profiles": {
"QuestionLibraryMQConsumer": {
"commandName": "Project",
"dotnetRunMessages": true,
"environmentVariables": {
"DOTNET_ENVIRONMENT": "Development"
}
}
}
}

View File

@ -0,0 +1,16 @@
<Project Sdk="Microsoft.NET.Sdk.Worker">
<PropertyGroup>
<TargetFramework>net8.0</TargetFramework>
<Nullable>enable</Nullable>
<ImplicitUsings>enable</ImplicitUsings>
<UserSecretsId>dotnet-QuestionLibraryMQConsumer-3847697e-8fac-49a3-9723-51eb362c06ff</UserSecretsId>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Microsoft.Extensions.Hosting" Version="8.0.1" />
<PackageReference Include="RabbitMQ.Client" Version="7.2.1" />
<PackageReference Include="Serilog.AspNetCore" Version="10.0.0" />
<PackageReference Include="SqlSugarCore" Version="5.1.4.214" />
</ItemGroup>
</Project>

View File

@ -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"
}
}
]
}
}

View File

@ -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"
}
}
]
}
}