添加 任务缓存 防止回调频繁修改数据

This commit is contained in:
2025-06-16 20:12:35 +08:00
parent bc90b17961
commit c12e1b5d65
13 changed files with 333 additions and 114 deletions
@@ -20,6 +20,8 @@ public static class QuartzTaskSchedulerConfig
// 每30秒任务配置
ConfigureThirtySecondTask(q, chinaTimeZone);
ConfigureFiftySecondTask(q, chinaTimeZone);
ConfigureFiveMinuteTask(q, chinaTimeZone);
});
@@ -33,6 +35,7 @@ public static class QuartzTaskSchedulerConfig
services.AddTransient<TokenResetService>();
services.AddTransient<TokenSyncService>();
services.AddTransient<TaskStatusCheckService>();
services.AddTransient<TaskSyncService>();
}
private static TimeZoneInfo GetChinaTimeZone()
@@ -88,6 +91,16 @@ public static class QuartzTaskSchedulerConfig
.WithCronSchedule("*/30 * * * * ?", x => x.InTimeZone(timeZone))); // 每30秒执行一次
}
private static void ConfigureFiftySecondTask(IServiceCollectionQuartzConfigurator q, TimeZoneInfo timeZone)
{
var jobKey = new JobKey("FiftySecondTask", "DefaultGroup");
q.AddJob<TaskSyncService>(opts => opts.WithIdentity(jobKey));
q.AddTrigger(opts => opts
.ForJob(jobKey)
.WithIdentity("FiftySecondTaskTrigger", "DefaultGroup")
.WithCronSchedule("*/15 * * * * ?", x => x.InTimeZone(timeZone))); // 每30秒执行一次
}
private static void ConfigureFiveMinuteTask(IServiceCollectionQuartzConfigurator q, TimeZoneInfo timeZone)
{
var jobKey = new JobKey("FiveMinuteTask", "DefaultGroup");
@@ -12,12 +12,13 @@ namespace LMS.service.Controllers
{
[Route("api/[controller]")]
[ApiController]
public class MJPackageController(TokenUsageTracker usageTracker, ITaskConcurrencyManager taskConcurrencyManager, ILogger<MJPackageController> logger, IMJPackageService mJPackageService) : ControllerBase
public class MJPackageController(TokenUsageTracker usageTracker, ITaskConcurrencyManager taskConcurrencyManager, ILogger<MJPackageController> logger, IMJPackageService mJPackageService, ITokenService tokenService) : ControllerBase
{
private readonly TokenUsageTracker _usageTracker = usageTracker;
private readonly ILogger<MJPackageController> _logger = logger;
private readonly ITaskConcurrencyManager _taskConcurrencyManager = taskConcurrencyManager;
private readonly IMJPackageService _mJPackageService = mJPackageService;
private readonly ITokenService _tokenService = tokenService;
[HttpPost("mj/submit/imagine")]
[RateLimit]
@@ -42,7 +43,10 @@ namespace LMS.service.Controllers
string body = JsonConvert.SerializeObject(model);
client.Timeout = Timeout.InfiniteTimeSpan;
string mjUrl = "https://api.laitool.cc/mj/submit/imagine";
string mjAPIBasicUrl = await _tokenService.GetMJAPIBasicUrl();
string mjUrl = $"{mjAPIBasicUrl}/mj/submit/imagine";
var response = await client.PostAsync(mjUrl, new StringContent(body, Encoding.UTF8, "application/json"));
// 读取响应内容
string content = await response.Content.ReadAsStringAsync();
@@ -150,11 +150,9 @@ namespace LMS.service.Extensions.Attributes
// 在异常情况下也要释放并发许可
if (_concurrencyAcquired)
{
usageTracker.ReleaseConcurrencyPermit(_token);
}
logger.LogError(ex, $"处理Token请求时发生错误: {_token},已释放Token许可!");
logger.LogError(ex, $"处理Token请求时发生错误: {_token}是否有并发 {_concurrencyAcquired}已释放Token许可!");
context.Result = new ObjectResult("Internal server error")
{
StatusCode = StatusCodes.Status500InternalServerError
@@ -6,7 +6,6 @@ using Microsoft.AspNetCore.Mvc;
using Microsoft.EntityFrameworkCore;
using System.Net.Sockets;
using System.Text.Json;
using static Betalgo.Ranul.OpenAI.ObjectModels.StaticValues.AssistantsStatics.MessageStatics;
namespace LMS.service.Service.MJPackage
{
@@ -94,6 +93,7 @@ namespace LMS.service.Service.MJPackage
{
TaskId = mjTask.TaskId,
Token = mjTask.Token,
TokenId = mjTask.TokenId,
Status = status,
StartTime = mjTask.StartTime,
EndTime = null,
@@ -106,20 +106,36 @@ namespace LMS.service.Service.MJPackage
// 当前任务已经被释放过了
// 开始修改数据
mJApiTasks.EndTime = BeijingTimeExtension.GetBeijingTime();
await _taskConcurrencyManager.UpdateTaskInDatabase(mJApiTasks);
// 不直接修改数据库了 改为修改缓存
bool modifySatus = _usageTracker.AddOrUpdateTaskCache(mJApiTasks);
if (modifySatus == false)
{
// 缓存修改失败,可能是因为任务不存在或状态不匹配,尝试修改数据库
await _taskConcurrencyManager.UpdateTaskInDatabase(mJApiTasks);
}
return new OkObjectResult(null);
}
if (status == MJTaskStatus.SUCCESS || status == MJTaskStatus.FAILURE || status == MJTaskStatus.CANCEL)
{
mJApiTasks.EndTime = BeijingTimeExtension.GetBeijingTime();
// 任务状态为成功、失败或取消,释放Token
_logger.LogInformation("MJNotifyHook 回调: 任务状态为成功、失败或取消,释放Token " + mjTask.Token);
// 开始修改数据,然后在释放
_usageTracker.ReleaseConcurrencyPermit(mjTask.Token);
}
else
{
// 只修改状态,不释放Token
mJApiTasks.EndTime = null; // 处理中没有结束时间
}
// 开始修改数据
await _taskConcurrencyManager.UpdateTaskInDatabase(mJApiTasks);
bool updateStatus = _usageTracker.AddOrUpdateTaskCache(mJApiTasks);
if (updateStatus == false)
{
// 缓存修改失败,可能是因为任务不存在或状态不匹配,尝试修改数据库
await _taskConcurrencyManager.UpdateTaskInDatabase(mJApiTasks);
}
return new OkObjectResult(null);
}
catch (Exception ex)
@@ -146,7 +162,7 @@ namespace LMS.service.Service.MJPackage
private async Task<ActionResult<object>?> TryOriginApiAsync(string id)
{
const string originUrl = "https://mjapi.bzu.cn/mj/task/{0}/fetch";
string originUrl = $"https://mjapi.bzu.cn/mj/task/{id}/fetch";
// 判断 原始token 不存在 直接 返回空
string orginToken = await _tokenService.GetOriginToken();
@@ -188,12 +204,12 @@ namespace LMS.service.Service.MJPackage
private async Task<ActionResult<object>?> TryBackupApiAsync(string id, string useToken)
{
const string backupUrlTemplate = "https://api.laitool.cc/mj/task/{0}/fetch";
string mjAPIBasicUrl = await _tokenService.GetMJAPIBasicUrl();
string backupUrl = $"{mjAPIBasicUrl}/mj/task/{id}/fetch";
const int maxRetries = 3;
const int baseDelayMs = 1000;
using var client = CreateHttpClient($"Bearer sk-{useToken}", true);
var backupUrl = string.Format(backupUrlTemplate, id);
for (int attempt = 1; attempt <= maxRetries; attempt++)
{