From 12d488a5ca18eab275d1d57ff98d357dd962d7d5 Mon Sep 17 00:00:00 2001
From: glz <694770232@qq.com>
Date: Thu, 2 Jul 2026 16:05:13 +0800
Subject: [PATCH] =?UTF-8?q?refactor:=20=E4=BC=98=E5=8C=96AI=E8=AF=84?=
=?UTF-8?q?=E5=88=86=E4=B8=8E=E4=BA=8C=E7=BB=B4=E7=A0=81ID=E7=94=9F?=
=?UTF-8?q?=E6=88=90=E9=80=BB=E8=BE=91=EF=BC=8C=E8=B0=83=E6=95=B4RabbitMQ?=
=?UTF-8?q?=E9=85=8D=E7=BD=AE?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
1. 调整RabbitMQ预取计数配置,支持从配置读取
2. 新增随机ID帮助类,生成唯一长整型ID
3. 重构二维码ID生成逻辑,新增重试机制避免重复
4. 优化AI评分配置,调整温度系数与并发限制
5. 重构跨页题评分逻辑,支持分组评分与结果去重
6. 新增AI评分异常分类与结果校验逻辑
7. 优化评分提示词与结果归一化处理
---
.../Helpers/RandomIdHelper.cs | 37 ++
.../RabbitMQ/RabbitMQService.cs | 13 +-
.../UserJournalService.cs | 54 +-
.../Consumers/JournalTaskReceiveConsumer.cs | 499 +++++++++++++++---
.../appsettings.json | 6 +-
5 files changed, 524 insertions(+), 85 deletions(-)
create mode 100644 QYZH.InteractiveMagazine.Common/Helpers/RandomIdHelper.cs
diff --git a/QYZH.InteractiveMagazine.Common/Helpers/RandomIdHelper.cs b/QYZH.InteractiveMagazine.Common/Helpers/RandomIdHelper.cs
new file mode 100644
index 0000000..78fd452
--- /dev/null
+++ b/QYZH.InteractiveMagazine.Common/Helpers/RandomIdHelper.cs
@@ -0,0 +1,37 @@
+using System.Security.Cryptography;
+
+namespace QYZH.InteractiveMagazine.Common.Helpers;
+
+///
+/// 随机ID帮助类
+///
+public static class RandomIdHelper
+{
+ private const long DefaultMinValue = 1_000_000_000_000_000_000L;
+ private const long DefaultMaxValue = 9_000_000_000_000_000_000L;
+
+ ///
+ /// 生成不可预测的正数长整型ID
+ ///
+ public static long GenerateLongId(long minValue = DefaultMinValue, long maxValue = DefaultMaxValue)
+ {
+ if (minValue <= 0 || minValue >= maxValue)
+ {
+ throw new ArgumentOutOfRangeException(nameof(minValue), "随机ID范围配置错误");
+ }
+
+ var range = (ulong)(maxValue - minValue);
+ var limit = ulong.MaxValue - (ulong.MaxValue % range);
+
+ Span bytes = stackalloc byte[sizeof(ulong)];
+ while (true)
+ {
+ RandomNumberGenerator.Fill(bytes);
+ var value = BitConverter.ToUInt64(bytes);
+ if (value < limit)
+ {
+ return minValue + (long)(value % range);
+ }
+ }
+ }
+}
diff --git a/QYZH.InteractiveMagazine.Infrastructure/RabbitMQ/RabbitMQService.cs b/QYZH.InteractiveMagazine.Infrastructure/RabbitMQ/RabbitMQService.cs
index 2a168e8..7d1cbc1 100644
--- a/QYZH.InteractiveMagazine.Infrastructure/RabbitMQ/RabbitMQService.cs
+++ b/QYZH.InteractiveMagazine.Infrastructure/RabbitMQ/RabbitMQService.cs
@@ -1,5 +1,6 @@
using RabbitMQ.Client;
using RabbitMQ.Client.Events;
+using Microsoft.Extensions.Configuration;
using System.Text;
using System.Text.Encodings.Web;
using System.Text.Json;
@@ -9,13 +10,15 @@ 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)
+ public RabbitMQService(IRabbitMQConnection connection, IConfiguration configuration)
{
_connection = connection ?? throw new ArgumentNullException(nameof(connection));
+ _configuration = configuration;
}
@@ -130,7 +133,13 @@ namespace QYZH.InteractiveMagazine.Infrastructure.RabbitMQ
public async Task ReceiveAsync(string exchange, string queueName, string routingKey, Func callback, CancellationToken cancellationToken = default)
{
var channel = await _connection.CreateChannel();
- //await channel.BasicQosAsync(0, 10, false); // 一次最多接收10条未确认的消息
+ 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);
diff --git a/QYZH.InteractiveMagazine.Service/UserJournalService.cs b/QYZH.InteractiveMagazine.Service/UserJournalService.cs
index 75f67e0..3f3b1b4 100644
--- a/QYZH.InteractiveMagazine.Service/UserJournalService.cs
+++ b/QYZH.InteractiveMagazine.Service/UserJournalService.cs
@@ -32,6 +32,7 @@ public class UserJournalService(
private const string QrCodeGenerateQueue = "mq.journal.qrcode.generate";
private const string QrCodeGenerateRoutingKey = "rk.journal.qrcode.generate";
private const int MaxBatchQrCodeCount = 500;
+ private const int MaxRandomIdGenerateRetryCount = 5;
///
/// 用户绑定期刊(扫码绑定)
@@ -183,8 +184,10 @@ public class UserJournalService(
throw new BusinessException("该期刊暂未发布,无法生成二维码", ResultCode.UNPROCESSABLE_ENTITY);
}
+ var recordId = await GenerateUniqueQrCodeIdAsync();
var record = new UserJournal
{
+ Id = recordId,
UserId = null,
JournalId = input.JournalId,
Type = 0,
@@ -243,9 +246,11 @@ public class UserJournalService(
}
var now = DateTime.Now;
- var records = Enumerable.Range(0, input.Count)
- .Select(_ => new UserJournal
+ var randomIds = await GenerateUniqueQrCodeIdsAsync(input.Count);
+ var records = randomIds
+ .Select(id => new UserJournal
{
+ Id = id,
UserId = null,
JournalId = input.JournalId,
Type = 0,
@@ -459,6 +464,51 @@ public class UserJournalService(
return JsonSerializer.Serialize(new { JournalId = journalId, Id = id });
}
+ private async Task GenerateUniqueQrCodeIdAsync()
+ {
+ for (var i = 0; i < MaxRandomIdGenerateRetryCount; i++)
+ {
+ var id = RandomIdHelper.GenerateLongId();
+ var exists = await userJournalRepository.Queryable()
+ .AnyAsync(uj => uj.Id == id);
+
+ if (!exists)
+ {
+ return id;
+ }
+ }
+
+ throw new BusinessException("生成二维码ID失败,请稍后重试", ResultCode.GLOBAL_ERROR);
+ }
+
+ private async Task> GenerateUniqueQrCodeIdsAsync(int count)
+ {
+ var ids = new HashSet();
+
+ for (var i = 0; i < MaxRandomIdGenerateRetryCount && ids.Count < count; i++)
+ {
+ while (ids.Count < count)
+ {
+ ids.Add(RandomIdHelper.GenerateLongId());
+ }
+
+ var candidateIds = ids.ToList();
+ var existingIds = await userJournalRepository.Queryable()
+ .Where(uj => candidateIds.Contains(uj.Id))
+ .Select(uj => uj.Id)
+ .ToListAsync();
+
+ if (existingIds.Count == 0)
+ {
+ return candidateIds;
+ }
+
+ ids.ExceptWith(existingIds);
+ }
+
+ throw new BusinessException("生成二维码ID失败,请稍后重试", ResultCode.GLOBAL_ERROR);
+ }
+
private async Task SendBindJournalMessageAsync(Users user, Journal journal)
{
try
diff --git a/QYZH.InteractiveMagazine.WorkService/Consumers/JournalTaskReceiveConsumer.cs b/QYZH.InteractiveMagazine.WorkService/Consumers/JournalTaskReceiveConsumer.cs
index 6152a7d..eac5b1b 100644
--- a/QYZH.InteractiveMagazine.WorkService/Consumers/JournalTaskReceiveConsumer.cs
+++ b/QYZH.InteractiveMagazine.WorkService/Consumers/JournalTaskReceiveConsumer.cs
@@ -21,9 +21,14 @@ public class JournalTaskReceiveConsumer(
{
private const int DefaultAiScoreMaxRetryCount = 3;
private const int DefaultAiScoreRetryDelayMilliseconds = 1000;
+ private const int DefaultAiMaxConcurrency = 1;
private const long DefaultMaxImageBytes = 10 * 1024 * 1024;
private const float DefaultCompletionThreshold = 80;
private const float DefaultCommunityScoreThreshold = 90;
+ private const string AiProcessingMessage = "AI批阅中,请稍后";
+ private static readonly object AiSemaphoreLock = new();
+ private static SemaphoreSlim? aiSemaphore;
+ private static int aiSemaphoreLimit;
public string Exchange => "ex.journal";
@@ -54,17 +59,38 @@ public class JournalTaskReceiveConsumer(
.Where(t => taskIds.Contains(t.Id) && !t.IsDeleted)
.ToListAsync(cancellationToken);
var taskMap = tasks.ToDictionary(t => t.Id);
+ var groupIds = tasks.Select(t => t.GroupId > 0 ? t.GroupId : t.Id).Distinct().ToList();
+ var groupTasks = await client.Queryable()
+ .Where(t => groupIds.Contains(t.GroupId) && !t.IsDeleted)
+ .ToListAsync(cancellationToken);
+ var persistedTaskMap = groupTasks
+ .GroupBy(t => t.GroupId > 0 ? t.GroupId : t.Id)
+ .ToDictionary(
+ g => g.Key,
+ g => g.FirstOrDefault(t => t.Id == g.Key)
+ ?? g.OrderBy(t => t.Id).First());
+ var groupTaskIds = groupTasks.Select(t => t.Id).Distinct().ToList();
var referenceAnswers = await client.Queryable()
- .Where(a => taskIds.Contains(a.JournalPageTaskId) && !a.IsDeleted)
+ .Where(a => groupTaskIds.Contains(a.JournalPageTaskId) && !a.IsDeleted)
.ToListAsync(cancellationToken);
var referenceAnswerMap = referenceAnswers
.GroupBy(a => a.JournalPageTaskId)
.ToDictionary(g => g.Key, g => g.ToList());
- var page = await client.Queryable()
- .Where(p => p.Id == data.PageId && !p.IsDeleted)
- .FirstAsync(cancellationToken);
+ var pageIds = groupTasks.Select(t => t.JournalPageId).Append(data.PageId).Distinct().ToList();
+ var pages = await client.Queryable()
+ .Where(p => pageIds.Contains(p.Id) && !p.IsDeleted)
+ .ToListAsync(cancellationToken);
+ var pageMap = pages.ToDictionary(p => p.Id);
+ var existingAnswers = await client.Queryable()
+ .Where(a => a.UserId == data.UserId && groupTaskIds.Contains(a.JournalPageTaskId) && !a.IsDeleted)
+ .ToListAsync(cancellationToken);
+ var existingAnswerMap = existingAnswers
+ .GroupBy(a => a.JournalPageTaskId)
+ .ToDictionary(g => g.Key, g => g.OrderByDescending(a => a.UpdatedAt ?? a.CreatedAt).First());
+
+ var questionContexts = new List();
var answerContexts = new List();
foreach (var question in data.Questions)
{
@@ -80,13 +106,55 @@ public class JournalTaskReceiveConsumer(
continue;
}
- referenceAnswerMap.TryGetValue(task.Id, out var taskReferenceAnswers);
- var scoreResult = await ScoreQuestionAsync(task, question, taskReferenceAnswers ?? [], cancellationToken);
- answerContexts.Add(new JournalAnswerContext(
- BuildAnswerEntity(data, question, task, page, scoreResult, GetCompletionThreshold()),
- question,
- task,
- scoreResult));
+ pageMap.TryGetValue(task.JournalPageId, out var taskPage);
+ questionContexts.Add(new JournalQuestionContext(question, task, taskPage));
+ }
+
+ foreach (var scoreUnit in BuildScoreUnits(questionContexts))
+ {
+ var persistedContext = BuildPersistedContext(scoreUnit, persistedTaskMap, pageMap);
+ if (persistedContext == null)
+ {
+ logger.LogWarning("期刊任务评分单元缺少可入库任务,UserId: {UserId}, GroupId: {GroupId}",
+ data.UserId, scoreUnit.GroupId);
+ continue;
+ }
+
+ if (IsCompletedSameAnswer(scoreUnit, persistedContext.Task, data, existingAnswerMap))
+ {
+ logger.LogInformation("期刊任务作答已完成且内容未变化,跳过AI评分,UserId: {UserId}, GroupId: {GroupId}, TaskIds: {TaskIds}",
+ data.UserId, scoreUnit.GroupId, string.Join(",", scoreUnit.Questions.Select(q => q.Task.Id)));
+ continue;
+ }
+
+ try
+ {
+ var scoreResult = await ScoreQuestionAsync(scoreUnit, referenceAnswerMap, cancellationToken);
+ var normalizedResult = NormalizeScoreResult(scoreResult, persistedContext.Task);
+ answerContexts.Add(new JournalAnswerContext(
+ BuildAnswerEntity(data, scoreUnit, persistedContext.Question, persistedContext.Task, persistedContext.Page, normalizedResult, GetCompletionThreshold()),
+ persistedContext.Question,
+ persistedContext.Task,
+ normalizedResult));
+ LogSkippedGroupTasks(scoreUnit, persistedContext.Task);
+ }
+ catch (OperationCanceledException)
+ {
+ throw;
+ }
+ catch (Exception ex)
+ {
+ logger.LogWarning(ex, "期刊任务AI评分失败,已标记为处理中,UserId: {UserId}, GroupId: {GroupId}, TaskIds: {TaskIds}",
+ data.UserId, scoreUnit.GroupId, string.Join(",", scoreUnit.Questions.Select(q => q.Task.Id)));
+
+ var failureResult = BuildFailureScoreResult();
+ answerContexts.Add(new JournalAnswerContext(
+ BuildAnswerEntity(data, scoreUnit, persistedContext.Question, persistedContext.Task, persistedContext.Page, failureResult, GetCompletionThreshold()),
+ persistedContext.Question,
+ persistedContext.Task,
+ failureResult));
+ LogSkippedGroupTasks(scoreUnit, persistedContext.Task);
+ }
}
if (answerContexts.Count == 0)
@@ -148,9 +216,8 @@ public class JournalTaskReceiveConsumer(
}
private async Task ScoreQuestionAsync(
- JournalPageTask task,
- Question question,
- List referenceAnswers,
+ JournalScoreUnit scoreUnit,
+ Dictionary> referenceAnswerMap,
CancellationToken cancellationToken)
{
var apiKey = configuration["AiChat:ApiKey"];
@@ -158,45 +225,61 @@ public class JournalTaskReceiveConsumer(
var model = configuration["AiChat:Model"];
var timeoutSeconds = configuration.GetValue("AiChat:TimeoutSeconds");
var maxTokens = configuration.GetValue("AiChat:MaxTokens");
- var temperature = configuration.GetValue("AiChat:Temperature");
+ var temperature = configuration.GetValue("AiChat:Temperature");
if (string.IsNullOrWhiteSpace(apiKey) || string.IsNullOrWhiteSpace(baseUrl) || string.IsNullOrWhiteSpace(model))
{
throw new InvalidOperationException("AI聊天服务配置不完整,请检查 AiChat 配置节点");
}
- var answerImages = await BuildAnswerImageContentsAsync(question, cancellationToken);
- if (answerImages.Count == 0)
- {
- throw new InvalidOperationException($"题目 {question.Id} 缺少答案图片");
- }
-
- var referenceAnswerImages = await BuildReferenceAnswerImageContentsAsync(task.Id, referenceAnswers, cancellationToken);
-
var content = new List