From 77efd0b2c794fffd27dae70227f2aa9bab6ec67b Mon Sep 17 00:00:00 2001
From: luogw <3132758203@qq.com>
Date: Wed, 15 Jul 2026 16:32:51 +0800
Subject: [PATCH] =?UTF-8?q?bug=E4=BF=AE=E5=A4=8D?=
MIME-Version: 1.0
Content-Type: text/plain; charset=UTF-8
Content-Transfer-Encoding: 8bit
---
...GenerateQuestionQueueSchedulerService.java | 19 +-
.../GenerateQuestionQueueServiceImpl.java | 293 ++++++++++++------
src/main/resources/application-dev.yml | 2 +-
3 files changed, 196 insertions(+), 118 deletions(-)
diff --git a/src/main/java/com/project/interaction/application/impl/GenerateQuestionQueueSchedulerService.java b/src/main/java/com/project/interaction/application/impl/GenerateQuestionQueueSchedulerService.java
index ca01b12..c22cd03 100644
--- a/src/main/java/com/project/interaction/application/impl/GenerateQuestionQueueSchedulerService.java
+++ b/src/main/java/com/project/interaction/application/impl/GenerateQuestionQueueSchedulerService.java
@@ -19,9 +19,9 @@ public class GenerateQuestionQueueSchedulerService {
/**
* 定时处理队列中的待重试项
- * 每分钟执行一次
+ * 上一次执行完成后间隔 60s 再次执行,避免重试耗时超过周期时背靠背/重入
*/
- @Scheduled(fixedRate = 60000)
+ @Scheduled(fixedDelay = 60000)
public void processQueuedItems() {
try {
log.debug(">>> [定时任务] 开始处理队列中的待重试题目生成请求...");
@@ -30,20 +30,5 @@ public class GenerateQuestionQueueSchedulerService {
log.error(">>> [定时任务] 处理队列时发生异常, 错误: {}", e.getMessage(), e);
}
}
-
- /**
- * 定时监控队列状态
- * 每5分钟执行一次,监控队列大小和消费状态
- */
- @Scheduled(fixedRate = 300000) // 每300秒执行一次(300000毫秒)
- public void monitorQueueStatus() {
- try {
- log.debug(">>> [定时任务] 开始监控队列状态...");
- // 这里可以添加队列监控逻辑,如记录队列大小、消费暂停/恢复状态等
- log.debug(">>> [定时任务] 队列监控完成");
- } catch (Exception e) {
- log.error(">>> [定时任务] 监控队列时发生异常, 错误: {}", e.getMessage(), e);
- }
- }
}
diff --git a/src/main/java/com/project/interaction/domain/service/impl/GenerateQuestionQueueServiceImpl.java b/src/main/java/com/project/interaction/domain/service/impl/GenerateQuestionQueueServiceImpl.java
index 8070565..3aad807 100644
--- a/src/main/java/com/project/interaction/domain/service/impl/GenerateQuestionQueueServiceImpl.java
+++ b/src/main/java/com/project/interaction/domain/service/impl/GenerateQuestionQueueServiceImpl.java
@@ -1,133 +1,194 @@
package com.project.interaction.domain.service.impl;
+import com.fasterxml.jackson.databind.DeserializationFeature;
+import com.fasterxml.jackson.databind.JsonNode;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.databind.SerializationFeature;
+import com.fasterxml.jackson.databind.json.JsonMapper;
+import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
import com.project.interaction.domain.dto.GenerateQuestionQueueDTO;
import com.project.interaction.domain.service.GenerateQuestionQueueService;
-import jakarta.annotation.PostConstruct;
import lombok.extern.slf4j.Slf4j;
+import org.redisson.api.RMap;
+import org.redisson.api.RScoredSortedSet;
+import org.redisson.api.RedissonClient;
+import org.redisson.client.codec.StringCodec;
+import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Service;
import java.time.LocalDateTime;
-import java.util.*;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.PriorityBlockingQueue;
-import java.util.concurrent.atomic.AtomicInteger;
-import java.util.concurrent.locks.ReentrantLock;
-
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.Comparator;
+import java.util.List;
+import java.util.Map;
+import java.util.UUID;
+
+/**
+ * 生成题目重试队列 - 基于 Redisson 的持久化实现
+ *
+ * 数据结构:
+ *
+ * - {@code items} (RMap) itemId -> DTO(JSON),队列内容的唯一真相来源
+ * - {@code schedule} (RScoredSortedSet) itemId,score = 下次可重试时间(epoch ms),用于延迟/退避
+ *
+ * 相比原内存实现,重启不丢失、支持指数退避;当前不限制重试次数,失败项将持续按退避重试。
+ */
@Service
@Slf4j
public class GenerateQuestionQueueServiceImpl implements GenerateQuestionQueueService {
+ /** 重试基础间隔(秒),退避从此值开始按 2 的幂增长 */
@Value("${question.queue.retry-interval:60}")
- private Long retryInterval;
-
- private PriorityBlockingQueue queue;
-
- /*按照权限降序、时间升序*/
- @PostConstruct
- public void init() {
- this.queue = new PriorityBlockingQueue<>(
- 10000,
- Comparator.comparing(GenerateQuestionQueueDTO::getWeight, Comparator.reverseOrder())
- .thenComparing(GenerateQuestionQueueDTO::getAddTime)
- .thenComparing(GenerateQuestionQueueDTO::getItemId)
- );
- }
+ private long retryIntervalSeconds;
+
+ /** 退避间隔上限(秒) */
+ @Value("${question.queue.max-retry-interval:1800}")
+ private long maxRetryIntervalSeconds;
+
+ private static final String ITEMS_KEY = "question:gen:queue:items";
+ private static final String SCHEDULE_KEY = "question:gen:queue:schedule";
- private final Map itemMap = new ConcurrentHashMap<>();
+ /** 自持 ObjectMapper:注册 JavaTime,日期写为 ISO 文本,忽略未知字段以保证前后兼容 */
+ private static final ObjectMapper MAPPER = JsonMapper.builder()
+ .addModule(new JavaTimeModule())
+ .disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS)
+ .configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false)
+ .build();
- private final AtomicInteger currentSize = new AtomicInteger(0);
+ @Autowired
+ private RedissonClient redissonClient;
- private final ReentrantLock capacityLock = new ReentrantLock();
+ private RMap items() {
+ return redissonClient.getMap(ITEMS_KEY, StringCodec.INSTANCE);
+ }
+
+ private RScoredSortedSet schedule() {
+ return redissonClient.getScoredSortedSet(SCHEDULE_KEY, StringCodec.INSTANCE);
+ }
@Override
public boolean addItem(GenerateQuestionQueueDTO item) {
- capacityLock.lock();
- try {
- if (item.getItemId() == null) {
- item.setItemId(UUID.randomUUID().toString());
- }
- if (item.getAddTime() == null) {
- item.setAddTime(LocalDateTime.now());
- }
- if (item.getRetryCount() == null) {
- item.setRetryCount(0);
- }
+ if (item == null) {
+ return false;
+ }
+ if (item.getItemId() == null) {
+ item.setItemId(UUID.randomUUID().toString());
+ }
+ if (item.getAddTime() == null) {
+ item.setAddTime(LocalDateTime.now());
+ }
+ if (item.getRetryCount() == null) {
+ item.setRetryCount(0);
+ }
- itemMap.put(item.getItemId(), item);
- queue.offer(item);
- currentSize.incrementAndGet();
+ String json = serialize(item);
+ if (json == null) {
+ return false;
+ }
- log.info(">>> [队列管理] 成功添加队列项, ItemId: {}, 当前队列大小: {}, 添加时间: {}",
- item.getItemId(), currentSize.get(), item.getAddTime());
+ RMap items = items();
+ // 先写内容,再登记调度:若中途异常,getRetryItems 会自愈清理孤立调度项
+ items.put(item.getItemId(), json);
+ schedule().add(System.currentTimeMillis(), item.getItemId());
- return true;
- } finally {
- capacityLock.unlock();
- }
+ log.info(">>> [队列管理] 成功添加队列项, ItemId: {}, 当前队列大小: {}, 添加时间: {}",
+ item.getItemId(), items.size(), item.getAddTime());
+ return true;
}
@Override
public List getRetryItems() {
- capacityLock.lock();
- try {
- List retryItems = new ArrayList<>();
-
- while (!queue.isEmpty()) {
- GenerateQuestionQueueDTO item = queue.poll();
- if (item == null) {
- break;
- }
-
- GenerateQuestionQueueDTO latestItem = itemMap.get(item.getItemId());
- if (latestItem == null) {
- // queue 中残留项(已被 removeItem/removeByTaskId 移除),直接丢弃
- continue;
- }
- retryItems.add(latestItem);
- }
+ long now = System.currentTimeMillis();
+ RMap items = items();
+ RScoredSortedSet schedule = schedule();
+
+ // 取出所有到期项(score <= now)。此处不移除,保证进程崩溃时重试项不丢失
+ Collection dueIds = schedule.valueRange(0, true, (double) now, true);
+ if (dueIds == null || dueIds.isEmpty()) {
+ log.debug(">>> [队列管理] 当前没有到期待重试的项");
+ return new ArrayList<>();
+ }
- return retryItems;
- } finally {
- capacityLock.unlock();
+ List retryItems = new ArrayList<>();
+ for (String itemId : dueIds) {
+ String json = items.get(itemId);
+ if (json == null) {
+ // 内容已被删除,清理残留调度项
+ schedule.remove(itemId);
+ continue;
+ }
+ GenerateQuestionQueueDTO item = deserialize(json);
+ if (item == null) {
+ // 反序列化失败(脏数据),丢弃避免卡死
+ items.remove(itemId);
+ schedule.remove(itemId);
+ continue;
+ }
+ retryItems.add(item);
}
+
+ // 保持原有优先级:权重降序、添加时间升序、itemId 升序
+ retryItems.sort(Comparator
+ .comparing(GenerateQuestionQueueDTO::getWeight,
+ Comparator.nullsLast(Comparator.reverseOrder()))
+ .thenComparing(GenerateQuestionQueueDTO::getAddTime,
+ Comparator.nullsLast(Comparator.naturalOrder()))
+ .thenComparing(GenerateQuestionQueueDTO::getItemId,
+ Comparator.nullsLast(Comparator.naturalOrder())));
+
+ log.info(">>> [队列管理] 获取到 {} 个待重试的项", retryItems.size());
+ return retryItems;
}
@Override
public boolean removeItem(String itemId) {
- capacityLock.lock();
- try {
- GenerateQuestionQueueDTO removed = itemMap.remove(itemId);
- if (removed != null) {
- currentSize.decrementAndGet();
- log.info(">>> [队列管理] 成功移除队列项, ItemId: {}, 当前队列大小: {}", itemId, currentSize.get());
- return true;
- }
+ if (itemId == null) {
return false;
- } finally {
- capacityLock.unlock();
}
+ boolean existed = items().remove(itemId) != null;
+ schedule().remove(itemId);
+ if (existed) {
+ log.info(">>> [队列管理] 成功移除队列项, ItemId: {}, 当前队列大小: {}", itemId, items().size());
+ }
+ return existed;
}
@Override
public int getQueueSize() {
- return currentSize.get();
+ return items().size();
}
+ @Override
public void requeue(GenerateQuestionQueueDTO item) {
if (item == null || item.getItemId() == null) {
return;
}
- capacityLock.lock();
- try {
- if (!itemMap.containsKey(item.getItemId())) {
- currentSize.incrementAndGet();
- }
- itemMap.put(item.getItemId(), item);
- queue.offer(item);
- } finally {
- capacityLock.unlock();
+ String itemId = item.getItemId();
+ RMap items = items();
+
+ // 仅当该项仍在追踪中时才重新入队;若已被 removeItem/removeByTaskId 删除,则不再复活
+ if (!items.containsKey(itemId)) {
+ log.info(">>> [队列管理] 项已被删除, 跳过重新入队, ItemId: {}", itemId);
+ return;
+ }
+
+ int retryCount = item.getRetryCount() == null ? 0 : item.getRetryCount();
+
+ // 指数退避:下次重试时间 = now + min(base * 2^(retryCount-1), cap);当前不限制重试次数
+ long delaySeconds = computeBackoffSeconds(retryCount);
+ long nextRetryMillis = System.currentTimeMillis() + delaySeconds * 1000L;
+
+ String json = serialize(item);
+ if (json == null) {
+ return;
}
+ items.put(itemId, json);
+ schedule().add((double) nextRetryMillis, itemId);
+
+ log.info(">>> [队列管理] 重新入队等待重试, ItemId: {}, 重试次数: {}, {}秒后重试",
+ itemId, retryCount, delaySeconds);
}
@Override
@@ -135,23 +196,55 @@ public class GenerateQuestionQueueServiceImpl implements GenerateQuestionQueueSe
if (taskId == null) {
return;
}
- capacityLock.lock();
- try {
- List keysToRemove = new ArrayList<>();
- for (Map.Entry entry : itemMap.entrySet()) {
- if (taskId.equals(entry.getValue().getTaskId())) {
- keysToRemove.add(entry.getKey());
- }
- }
- for (String key : keysToRemove) {
- itemMap.remove(key);
- currentSize.decrementAndGet();
- }
- if (!keysToRemove.isEmpty()) {
- log.info(">>> [队列管理] 根据TaskId移除队列项, TaskId: {}, 移除数量: {}", taskId, keysToRemove.size());
+ RMap items = items();
+ RScoredSortedSet schedule = schedule();
+
+ List keysToRemove = new ArrayList<>();
+ for (Map.Entry entry : items.readAllEntrySet()) {
+ if (taskId.equals(parseTaskId(entry.getValue()))) {
+ keysToRemove.add(entry.getKey());
}
- } finally {
- capacityLock.unlock();
+ }
+ for (String key : keysToRemove) {
+ items.remove(key);
+ schedule.remove(key);
+ }
+ if (!keysToRemove.isEmpty()) {
+ log.info(">>> [队列管理] 根据TaskId移除队列项, TaskId: {}, 移除数量: {}", taskId, keysToRemove.size());
+ }
+ }
+
+ /** 计算退避秒数:base * 2^(retryCount-1),上限 maxRetryIntervalSeconds */
+ private long computeBackoffSeconds(int retryCount) {
+ int exp = Math.min(Math.max(0, retryCount - 1), 20);
+ long delay = (long) (retryIntervalSeconds * Math.pow(2, exp));
+ return Math.min(delay, maxRetryIntervalSeconds);
+ }
+
+ private String serialize(GenerateQuestionQueueDTO item) {
+ try {
+ return MAPPER.writeValueAsString(item);
+ } catch (Exception e) {
+ log.error(">>> [队列管理] 序列化队列项失败, ItemId: {}, 原因: {}", item.getItemId(), e.getMessage(), e);
+ return null;
+ }
+ }
+
+ private GenerateQuestionQueueDTO deserialize(String json) {
+ try {
+ return MAPPER.readValue(json, GenerateQuestionQueueDTO.class);
+ } catch (Exception e) {
+ log.error(">>> [队列管理] 反序列化队列项失败, 原因: {}, 内容: {}", e.getMessage(), json);
+ return null;
+ }
+ }
+
+ private Long parseTaskId(String json) {
+ try {
+ JsonNode node = MAPPER.readTree(json).get("taskId");
+ return node == null || node.isNull() ? null : node.asLong();
+ } catch (Exception e) {
+ return null;
}
}
}
diff --git a/src/main/resources/application-dev.yml b/src/main/resources/application-dev.yml
index 12808d5..732db32 100644
--- a/src/main/resources/application-dev.yml
+++ b/src/main/resources/application-dev.yml
@@ -95,7 +95,7 @@ question:
generation:
# 限流速率:每秒允许的API请求数
rate-limit: 20
- downgrade: true
+ downgrade: false
# 队列配置
queue:
# 重试间隔(秒)