diff --git a/QYZH.InteractiveMagazine.Models/Settings/HangfireJobSettings.cs b/QYZH.InteractiveMagazine.Models/Settings/HangfireJobSettings.cs new file mode 100644 index 0000000..f7cd734 --- /dev/null +++ b/QYZH.InteractiveMagazine.Models/Settings/HangfireJobSettings.cs @@ -0,0 +1,56 @@ +namespace QYZH.InteractiveMagazine.Models.Settings; + +/// +/// Hangfire 定时任务配置集合(对应 appsettings.json 中的 "HangfireJobs" 节点) +/// +public class HangfireJobSettings +{ + /// + /// 定时任务列表,每一项对应一个 Hangfire RecurringJob + /// + public List Jobs { get; set; } = new(); +} + +/// +/// 单个定时任务配置 +/// +public class HangfireJobConfig +{ + /// + /// 任务名称(Hangfire Dashboard 中展示的唯一标识,建议英文短横线命名,如 "sample-job") + /// + public string Name { get; set; } = ""; + + /// + /// Job 类的完整类型名称(含命名空间,如 "QYZH.InteractiveMagazine.WorkService.Jobs.SampleJob") + /// + public string JobType { get; set; } = ""; + + /// + /// 要调用的方法名称(默认 "ExecuteAsync",方法必须为 public 无参 Task) + /// + public string MethodName { get; set; } = "ExecuteAsync"; + + /// + /// Cron 表达式,控制执行频率。常用示例: + /// "* * * * *" 每分钟 + /// "*/5 * * * *" 每5分钟 + /// "0 * * * *" 每小时整点 + /// "0 0 * * *" 每天午夜 + /// "0 9 * * *" 每天上午9点 + /// "0 9 * * 1" 每周一上午9点 + /// "0 0 1 * *" 每月1号午夜 + /// 格式:分 时 日 月 星期(5位,不支持秒和年) + /// + public string Cron { get; set; } = ""; + + /// + /// 是否启用此任务(设为 false 将从 Hangfire 中移除该任务) + /// + public bool Enabled { get; set; } = true; + + /// + /// 备注说明(仅供阅读理解,不参与运行逻辑) + /// + public string? Description { get; set; } +} diff --git a/QYZH.InteractiveMagazine.Service/PetService.cs b/QYZH.InteractiveMagazine.Service/PetService.cs index bc6f01c..07c3b41 100644 --- a/QYZH.InteractiveMagazine.Service/PetService.cs +++ b/QYZH.InteractiveMagazine.Service/PetService.cs @@ -252,7 +252,7 @@ public class PetService( logger.LogWarning("喂养失败,宠物未激活,PetId: {PetId}, Status: {Status}", input.PetId, pet.Status); throw new BusinessException("宠物未激活,无法喂养", 400); } - + FeedPetOutput result = new FeedPetOutput (); // 事务保证一致性 await UseTranAsync(async () => { diff --git a/QYZH.InteractiveMagazine.WebApi/Controllers/JournalController.cs b/QYZH.InteractiveMagazine.WebApi/Controllers/JournalController.cs index c37c484..6bceea4 100644 --- a/QYZH.InteractiveMagazine.WebApi/Controllers/JournalController.cs +++ b/QYZH.InteractiveMagazine.WebApi/Controllers/JournalController.cs @@ -22,7 +22,7 @@ namespace QYZH.InteractiveMagazine.WebApi.Controllers /// 期刊管理 /// [ApiExplorerSettings(GroupName = nameof(ApiVersionEnum.Platform))] - [Route("api/icr")] + [Route("api")] public class JournalController( OssService ossService, IJournalCatalogService JournalCatalogService, diff --git a/QYZH.InteractiveMagazine.WorkService/Consumers/IQueueConsumer.cs b/QYZH.InteractiveMagazine.WorkService/Consumers/IQueueConsumer.cs new file mode 100644 index 0000000..002ff9d --- /dev/null +++ b/QYZH.InteractiveMagazine.WorkService/Consumers/IQueueConsumer.cs @@ -0,0 +1,22 @@ +namespace QYZH.InteractiveMagazine.WorkService.Consumers; + +/// +/// 队列消费者接口 +/// +public interface IQueueConsumer +{ + /// + /// 监听的队列名称 + /// + string QueueName { get; } + + /// + /// 处理消息 + /// + Task HandleAsync(byte[] message); + + /// + /// 处理消费异常 + /// + Task OnErrorAsync(byte[] message, Exception exception); +} diff --git a/QYZH.InteractiveMagazine.WorkService/Consumers/JournalTaskReceiveConsumer.cs b/QYZH.InteractiveMagazine.WorkService/Consumers/JournalTaskReceiveConsumer.cs new file mode 100644 index 0000000..af30191 --- /dev/null +++ b/QYZH.InteractiveMagazine.WorkService/Consumers/JournalTaskReceiveConsumer.cs @@ -0,0 +1,35 @@ +using System.Text; + +namespace QYZH.InteractiveMagazine.WorkService.Consumers; + +/// +/// 期刊任务接收消费者(示例) +/// +public class JournalTaskReceiveConsumer : IQueueConsumer +{ + private readonly ILogger _logger; + + public JournalTaskReceiveConsumer(ILogger logger) + { + _logger = logger; + } + + public string QueueName => "mq.Journal.task.receive"; + + public async Task HandleAsync(byte[] message) + { + var body = Encoding.UTF8.GetString(message); + _logger.LogInformation("收到期刊任务消息: {Message}", body); + + // TODO: 在此编写具体的消息处理逻辑 + + await Task.CompletedTask; + } + + public Task OnErrorAsync(byte[] message, Exception exception) + { + var body = Encoding.UTF8.GetString(message); + _logger.LogError(exception, "处理期刊任务消息失败: {Message}", body); + return Task.CompletedTask; + } +} diff --git a/QYZH.InteractiveMagazine.WorkService/Consumers/RabbitMQHostedService.cs b/QYZH.InteractiveMagazine.WorkService/Consumers/RabbitMQHostedService.cs new file mode 100644 index 0000000..dee36af --- /dev/null +++ b/QYZH.InteractiveMagazine.WorkService/Consumers/RabbitMQHostedService.cs @@ -0,0 +1,82 @@ +using QYZH.InteractiveMagazine.Infrastructure.RabbitMQ; + +namespace QYZH.InteractiveMagazine.WorkService.Consumers; + +/// +/// RabbitMQ 消费者后台服务,自动发现并启动所有已注册的 IQueueConsumer +/// +public class RabbitMQHostedService : BackgroundService +{ + private readonly IServiceProvider _serviceProvider; + private readonly ILogger _logger; + + public RabbitMQHostedService(IServiceProvider serviceProvider, ILogger logger) + { + _serviceProvider = serviceProvider; + _logger = logger; + } + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + using var scope = _serviceProvider.CreateScope(); + var consumers = scope.ServiceProvider.GetServices().ToList(); + + if (consumers.Count == 0) + { + _logger.LogWarning("未注册任何队列消费者"); + return; + } + + _logger.LogInformation("发现 {Count} 个队列消费者,开始启动...", consumers.Count); + + var tasks = consumers.Select(c => StartConsumerAsync(c.QueueName, stoppingToken)).ToList(); + await Task.WhenAll(tasks); + } + + private async Task StartConsumerAsync(string queueName, CancellationToken stoppingToken) + { + _logger.LogInformation("正在启动消费者: {QueueName}", queueName); + + try + { + var rabbitMQService = _serviceProvider.GetRequiredService(); + + await rabbitMQService.ReceiveAsync(queueName, async (channel, ea) => + { + // 每条消息创建独立 scope,确保消费者内可注入 Scoped 服务 + using var messageScope = _serviceProvider.CreateScope(); + var consumer = messageScope.ServiceProvider + .GetServices() + .First(c => c.QueueName == queueName); + + try + { + await consumer.HandleAsync(ea.Body.ToArray()); + await channel.BasicAckAsync(ea.DeliveryTag, false, stoppingToken); + } + catch (Exception ex) + { + _logger.LogError(ex, "消费者 {QueueName} 处理消息异常", queueName); + await channel.BasicNackAsync(ea.DeliveryTag, false, true, stoppingToken); + + try + { + await consumer.OnErrorAsync(ea.Body.ToArray(), ex); + } + catch (Exception errorEx) + { + _logger.LogError(errorEx, "消费者 {QueueName} OnErrorAsync 执行异常", queueName); + } + } + }, stoppingToken); + } + catch (OperationCanceledException) + { + _logger.LogInformation("消费者 {QueueName} 已停止", queueName); + } + catch (Exception ex) + { + _logger.LogError(ex, "消费者 {QueueName} 启动失败", queueName); + } + } +} diff --git a/QYZH.InteractiveMagazine.WorkService/Jobs/SampleJob.cs b/QYZH.InteractiveMagazine.WorkService/Jobs/SampleJob.cs new file mode 100644 index 0000000..39c8912 --- /dev/null +++ b/QYZH.InteractiveMagazine.WorkService/Jobs/SampleJob.cs @@ -0,0 +1,30 @@ +using Microsoft.Extensions.Logging; + +namespace QYZH.InteractiveMagazine.WorkService.Jobs; + +/// +/// 示例定时任务 +/// +public class SampleJob +{ + private readonly ILogger _logger; + + public SampleJob(ILogger logger) + { + _logger = logger; + } + + /// + /// 执行定时任务(由Hangfire调度) + /// + public async Task ExecuteAsync() + { + _logger.LogInformation("SampleJob 开始执行 {Time}", DateTime.Now.ToString("yyyy-MM-dd HH:mm:ss")); + + // TODO: 在此编写具体的定时任务逻辑 + + await Task.CompletedTask; + + _logger.LogInformation("SampleJob 执行完毕"); + } +} diff --git a/QYZH.InteractiveMagazine.WorkService/Program.cs b/QYZH.InteractiveMagazine.WorkService/Program.cs new file mode 100644 index 0000000..44d3285 --- /dev/null +++ b/QYZH.InteractiveMagazine.WorkService/Program.cs @@ -0,0 +1,120 @@ +using Hangfire; +using Hangfire.Dashboard; +using Hangfire.MemoryStorage; +using QYZH.InteractiveMagazine.Infrastructure.RabbitMQ; +using QYZH.InteractiveMagazine.Models.Settings; +using QYZH.InteractiveMagazine.WorkService.Consumers; +using QYZH.InteractiveMagazine.WorkService.Jobs; +using Serilog; +using SqlSugar; +using SqlSugar.IOC; +using System.Linq.Expressions; +using System.Reflection; +using Yitter.IdGenerator; + +var builder = WebApplication.CreateBuilder(args); + +// 加载配置 +builder.Configuration + .AddJsonFile("appsettings.json", optional: false, reloadOnChange: true) + .AddJsonFile($"appsettings.{builder.Environment.EnvironmentName}.json", optional: true, reloadOnChange: true); + +// 配置Serilog +Log.Logger = new LoggerConfiguration() + .ReadFrom.Configuration(builder.Configuration) + .Enrich.FromLogContext() + .CreateLogger(); +builder.Services.AddSerilog(); + +// 初始化雪花ID生成器 +YitIdHelper.SetIdGenerator(new IdGeneratorOptions { WorkerId = 2 }); + +// 初始化MySQL(SqlSugar) +builder.Services.AddSqlSugar(new IocConfig +{ + ConfigId = 0, + DbType = IocDbType.MySql, + ConnectionString = builder.Configuration.GetConnectionString("DefaultConnection"), + IsAutoCloseConnection = true +}); +SugarIocServices.ConfigurationSugar(db => +{ + db.Aop.OnLogExecuting = (sql, pars) => + { + Log.Information("[SQL] {Sql}", UtilMethods.GetSqlString((DbType)IocDbType.MySql, sql, pars)); + }; + db.Aop.OnError = ex => + { + Log.Error(ex, "[SQL Error] {Message}", ex.Message); + }; +}); + +// 配置Hangfire(内存存储,后续可切换Redis) +builder.Services.AddHangfire(config => config + .UseMemoryStorage() + .UseSerializerSettings(new Newtonsoft.Json.JsonSerializerSettings + { + TypeNameHandling = Newtonsoft.Json.TypeNameHandling.All + })); +builder.Services.AddHangfireServer(); + +// 配置RabbitMQ +builder.Services.AddRabbitMQ(builder.Configuration); + +// 注册队列消费者(新增消费者只需实现 IQueueConsumer 并在此注册) +builder.Services.AddScoped(); + +// 注册消费者后台服务 +builder.Services.AddHostedService(); + +// 从配置文件读取定时任务列表 +var jobSettings = builder.Configuration.GetSection("HangfireJobs").Get(); + +var app = builder.Build(); + +// 配置Hangfire Dashboard(仅本机访问) +app.UseHangfireDashboard("/hangfire", new DashboardOptions +{ + Authorization = new[] { new LocalRequestsOnlyAuthorizationFilter() } +}); + +// 根据配置动态注册定时任务 +if (jobSettings?.Jobs != null) +{ + foreach (var job in jobSettings.Jobs) + { + var jobType = Type.GetType(job.JobType); + if (jobType == null) + { + Log.Warning("定时任务 [{Name}] 类型未找到: {JobType},跳过注册", job.Name, job.JobType); + continue; + } + + if (!job.Enabled) + { + // 配置为关闭的任务,从 Hangfire 中移除 + RecurringJob.RemoveIfExists(job.Name); + Log.Information("定时任务 [{Name}] 已禁用,已移除", job.Name); + continue; + } + + // 构造表达式:job => job.MethodName() + var method = jobType.GetMethod(job.MethodName); + if (method == null) + { + Log.Warning("定时任务 [{Name}] 方法未找到: {MethodName},跳过注册", job.Name, job.MethodName); + continue; + } + + var param = Expression.Parameter(jobType, "job"); + var call = Expression.Call(param, method); + var lambda = Expression.Lambda(call, param); + + RecurringJob.AddOrUpdate(job.Name, lambda, job.Cron); + Log.Information("定时任务 [{Name}] 已注册,Cron: {Cron}", job.Name, job.Cron); + } +} + +Log.Information("WorkService 已启动,Hangfire Dashboard: /hangfire"); + +app.Run(); diff --git a/QYZH.InteractiveMagazine.WorkService/Properties/launchSettings.json b/QYZH.InteractiveMagazine.WorkService/Properties/launchSettings.json new file mode 100644 index 0000000..58bc520 --- /dev/null +++ b/QYZH.InteractiveMagazine.WorkService/Properties/launchSettings.json @@ -0,0 +1,12 @@ +{ + "profiles": { + "QYZH.InteractiveMagazine.WorkService": { + "commandName": "Project", + "launchBrowser": true, + "environmentVariables": { + "ASPNETCORE_ENVIRONMENT": "Development" + }, + "applicationUrl": "https://localhost:53991;http://localhost:53992" + } + } +} \ No newline at end of file diff --git a/QYZH.InteractiveMagazine.WorkService/QYZH.InteractiveMagazine.WorkService.csproj b/QYZH.InteractiveMagazine.WorkService/QYZH.InteractiveMagazine.WorkService.csproj new file mode 100644 index 0000000..b931f2a --- /dev/null +++ b/QYZH.InteractiveMagazine.WorkService/QYZH.InteractiveMagazine.WorkService.csproj @@ -0,0 +1,25 @@ + + + + net8.0 + enable + enable + + + + + + + + + + + + + + + + + + + diff --git a/QYZH.InteractiveMagazine.WorkService/appsettings.json b/QYZH.InteractiveMagazine.WorkService/appsettings.json new file mode 100644 index 0000000..ca3a55b --- /dev/null +++ b/QYZH.InteractiveMagazine.WorkService/appsettings.json @@ -0,0 +1,48 @@ +{ + "ConnectionStrings": { + "DefaultConnection": "server=192.168.20.150;port=13306;database=InteractiveMagazine;user=user;password=n68792bu!y99r905;charset=utf8mb4;" + }, + "RabbitMq": { + "HostName": "192.168.20.150", + "Port": 5672, + "UserName": "smartschool", + "Password": "@ss%&*otz%d*pq2S", + "VirtualHost": "InteractiveMagazine" + }, + "Serilog": { + "MinimumLevel": { + "Default": "Information", + "Override": { + "Microsoft": "Warning", + "System": "Warning", + "Hangfire": "Information" + } + }, + "WriteTo": [ + { + "Name": "Console" + }, + { + "Name": "File", + "Args": { + "path": "logs/worker-log-.txt", + "rollingInterval": "Day" + } + } + ] + }, + "HangfireJobs": { + "_说明": "定时任务配置。新增任务只需在 Jobs 数组中追加一项。Cron格式:分 时 日 月 星期(5位)。常用:* * * * * 每分钟 | */5 * * * * 每5分钟 | 0 * * * * 每小时 | 0 9 * * * 每天9点 | 0 0 * * 1 每周一午夜", + "Jobs": [ + { + "Name": "sample-job", + "JobType": "QYZH.InteractiveMagazine.WorkService.Jobs.SampleJob", + "MethodName": "ExecuteAsync", + "Cron": "* * * * *", + "Enabled": false, + "Description": "示例任务(默认关闭,仅用于验证框架运行)" + } + ] + }, + "AllowedHosts": "*" +} diff --git a/QYZH.InteractiveMagazine.slnx b/QYZH.InteractiveMagazine.slnx index 76e4735..ca05cf1 100644 --- a/QYZH.InteractiveMagazine.slnx +++ b/QYZH.InteractiveMagazine.slnx @@ -6,4 +6,5 @@ +