diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/consumer/TimeRangeReConsumer.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/consumer/TimeRangeReConsumer.java index b9937f1..8748751 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/consumer/TimeRangeReConsumer.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/consumer/TimeRangeReConsumer.java @@ -17,9 +17,9 @@ import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus; import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently; import org.apache.rocketmq.common.consumer.ConsumeFromWhere; import org.apache.rocketmq.common.message.MessageExt; -import org.apache.rocketmq.common.message.MessageQueue; import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; import java.nio.charset.StandardCharsets; @@ -37,6 +37,9 @@ public class TimeRangeReConsumer implements DisposableBean { @Autowired(required = false) private MessageHistoryDataHandler messageHandler; + // 从配置文件读取队列数量 + @Value("${queues:4}") + private int totalQueues; // 测试环境配置为4,正式环境配置为16 // 时间范围配置 private static final String START_TIME = "2026-02-02 06:00:00"; @@ -184,16 +187,29 @@ public class TimeRangeReConsumer implements DisposableBean { } } + // 2. 超过结束时间,标记队列完成 // 2. 超过结束时间,标记队列完成 if (bornTimestamp > endTimestamp) { - log.info("队列 {} 已超过结束时间,标记为完成", queueId); - completedQueues.add(queueId); + boolean allQueuesNowCompleted = false; - // 检查是否所有队列都完成 - if (completedQueues.size() >= 4) { - log.info("所有队列均已超过结束时间,第一阶段完成"); - shouldContinue.set(false); + synchronized (completedQueues) { + boolean isNewCompletion = completedQueues.add(queueId); + if (isNewCompletion) { + log.info("队列 {} 已超过结束时间,标记为完成", queueId); + } + + // 检查是否所有队列都完成 - 必须在同步块中检查 + if (completedQueues.size() >= totalQueues && shouldContinue.get()) { + allQueuesNowCompleted = true; + shouldContinue.set(false); + } } + + // 在同步块外打印完成日志,确保只打印一次 + if (allQueuesNowCompleted) { + log.info("所有队列均已超过结束时间,第一阶段完成"); + } + return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } @@ -226,9 +242,6 @@ public class TimeRangeReConsumer implements DisposableBean { consumer.start(); log.info("时间范围重消费者启动成功,消费者组: {}", consumerGroup); - // 预热:查找大致的起始位置 - findApproximateStartOffset(consumer, startTimestamp); - // 监控第一阶段完成 waitForPhase1Completion(consumer); @@ -329,8 +342,8 @@ public class TimeRangeReConsumer implements DisposableBean { // 每30秒打印一次进度 if (checkCount % 6 == 0) { - log.info("第一阶段进度: 已接收={}, 已缓存={}, 设备数={}, 完成队列={}/4", - currentMessageCount, totalCached, messageCache.size(), completedQueues.size()); + log.info("第一阶段进度: 已接收={}, 已缓存={}, 设备数={}, 完成队列={}/{}", + currentMessageCount, totalCached, messageCache.size(), completedQueues.size(), totalQueues); } // 检查是否有进度 @@ -342,7 +355,7 @@ public class TimeRangeReConsumer implements DisposableBean { } // 如果30秒内没有新消息,认为第一阶段完成 - if (noProgressCount >= 6 && completedQueues.size() >= 4) { + if (noProgressCount >= 6 && completedQueues.size() >= totalQueues) { log.info("30秒内无新消息且所有队列完成,第一阶段完成"); shouldContinue.set(false); break; @@ -542,65 +555,6 @@ public class TimeRangeReConsumer implements DisposableBean { log.info("✓ 重消费者状态已重置"); } - /** - * 查找大致的起始offset - 修复版(处理异常) - */ - private void findApproximateStartOffset(DefaultMQPushConsumer consumer, long startTimestamp) throws Exception { - SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss"); - log.info("查找起始offset,目标时间: {}", sdf.format(new Date(startTimestamp))); - - try { - // 获取主题队列 - Set queues = consumer.fetchSubscribeMessageQueues(MessageCodeEnum.TR_AGENT_UP.getCode()); - - if (queues == null || queues.isEmpty()) { - log.info("无队列信息,跳过offset查找"); - return; - } - - log.info("找到{}个队列,开始设置起始offset...", queues.size()); - - for (MessageQueue queue : queues) { - try { - // 搜索指定时间戳的offset - long offset = consumer.searchOffset(queue, startTimestamp); - - if (offset >= 0) { - // 获取队列的offset范围 - long minOffset = consumer.minOffset(queue); - long maxOffset = consumer.maxOffset(queue); - - // 确保offset在有效范围内 - if (offset < minOffset) offset = minOffset; - if (offset > maxOffset) offset = maxOffset; - - queueCurrentOffsets.put(queue.getQueueId(), offset); - log.info("队列{}: offset={} (范围: {}-{})", - queue.getQueueId(), offset, minOffset, maxOffset); - } else { - // 使用最小offset - 使用新方法避免过时API - try { - long minOffset = consumer.minOffset(queue); - queueCurrentOffsets.put(queue.getQueueId(), minOffset); - log.info("队列{}: 使用最小offset={}", queue.getQueueId(), minOffset); - } catch (Exception e) { - // 如果minOffset也失败,使用0作为兜底 - log.warn("队列{}获取最小offset失败,使用0作为兜底: {}", queue.getQueueId(), e.getMessage()); - queueCurrentOffsets.put(queue.getQueueId(), 0L); - } - } - } catch (Exception e) { - log.warn("队列{}设置offset失败,使用0作为兜底: {}", queue.getQueueId(), e.getMessage()); - queueCurrentOffsets.put(queue.getQueueId(), 0L); - } - } - - log.info("起始offset设置完成"); - - } catch (Exception e) { - log.warn("查找起始offset失败,将使用默认时间戳方式: {}", e.getMessage()); - } - } /** * 打印第一阶段进度 @@ -611,8 +565,8 @@ public class TimeRangeReConsumer implements DisposableBean { totalCached += queue.size(); } - log.info("第一阶段进度: 已接收={}, 已缓存={}, 设备数={}, 完成队列={}/4", - receivedCount.get(), totalCached, messageCache.size(), completedQueues.size()); + log.info("第一阶段进度: 已接收={}, 已缓存={}, 设备数={}, 完成队列={}/{}", + receivedCount.get(), totalCached, messageCache.size(), completedQueues.size(), totalQueues); } /**