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 的持久化实现 + *

+ * 数据结构: + *

+ * 相比原内存实现,重启不丢失、支持指数退避;当前不限制重试次数,失败项将持续按退避重试。 + */ @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: # 重试间隔(秒)