优化历史消息处理
This commit is contained in:
+28
-74
@@ -17,9 +17,9 @@ import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
|
|||||||
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
|
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
|
||||||
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
|
import org.apache.rocketmq.common.consumer.ConsumeFromWhere;
|
||||||
import org.apache.rocketmq.common.message.MessageExt;
|
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.DisposableBean;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
|
import org.springframework.beans.factory.annotation.Value;
|
||||||
import org.springframework.stereotype.Component;
|
import org.springframework.stereotype.Component;
|
||||||
|
|
||||||
import java.nio.charset.StandardCharsets;
|
import java.nio.charset.StandardCharsets;
|
||||||
@@ -37,6 +37,9 @@ public class TimeRangeReConsumer implements DisposableBean {
|
|||||||
|
|
||||||
@Autowired(required = false)
|
@Autowired(required = false)
|
||||||
private MessageHistoryDataHandler messageHandler;
|
private MessageHistoryDataHandler messageHandler;
|
||||||
|
// 从配置文件读取队列数量
|
||||||
|
@Value("${queues:4}")
|
||||||
|
private int totalQueues; // 测试环境配置为4,正式环境配置为16
|
||||||
|
|
||||||
// 时间范围配置
|
// 时间范围配置
|
||||||
private static final String START_TIME = "2026-02-02 06:00:00";
|
private static final String START_TIME = "2026-02-02 06:00:00";
|
||||||
@@ -184,16 +187,29 @@ public class TimeRangeReConsumer implements DisposableBean {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 2. 超过结束时间,标记队列完成
|
||||||
// 2. 超过结束时间,标记队列完成
|
// 2. 超过结束时间,标记队列完成
|
||||||
if (bornTimestamp > endTimestamp) {
|
if (bornTimestamp > endTimestamp) {
|
||||||
log.info("队列 {} 已超过结束时间,标记为完成", queueId);
|
boolean allQueuesNowCompleted = false;
|
||||||
completedQueues.add(queueId);
|
|
||||||
|
|
||||||
// 检查是否所有队列都完成
|
synchronized (completedQueues) {
|
||||||
if (completedQueues.size() >= 4) {
|
boolean isNewCompletion = completedQueues.add(queueId);
|
||||||
log.info("所有队列均已超过结束时间,第一阶段完成");
|
if (isNewCompletion) {
|
||||||
shouldContinue.set(false);
|
log.info("队列 {} 已超过结束时间,标记为完成", queueId);
|
||||||
|
}
|
||||||
|
|
||||||
|
// 检查是否所有队列都完成 - 必须在同步块中检查
|
||||||
|
if (completedQueues.size() >= totalQueues && shouldContinue.get()) {
|
||||||
|
allQueuesNowCompleted = true;
|
||||||
|
shouldContinue.set(false);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// 在同步块外打印完成日志,确保只打印一次
|
||||||
|
if (allQueuesNowCompleted) {
|
||||||
|
log.info("所有队列均已超过结束时间,第一阶段完成");
|
||||||
|
}
|
||||||
|
|
||||||
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
|
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -226,9 +242,6 @@ public class TimeRangeReConsumer implements DisposableBean {
|
|||||||
consumer.start();
|
consumer.start();
|
||||||
log.info("时间范围重消费者启动成功,消费者组: {}", consumerGroup);
|
log.info("时间范围重消费者启动成功,消费者组: {}", consumerGroup);
|
||||||
|
|
||||||
// 预热:查找大致的起始位置
|
|
||||||
findApproximateStartOffset(consumer, startTimestamp);
|
|
||||||
|
|
||||||
// 监控第一阶段完成
|
// 监控第一阶段完成
|
||||||
waitForPhase1Completion(consumer);
|
waitForPhase1Completion(consumer);
|
||||||
|
|
||||||
@@ -329,8 +342,8 @@ public class TimeRangeReConsumer implements DisposableBean {
|
|||||||
|
|
||||||
// 每30秒打印一次进度
|
// 每30秒打印一次进度
|
||||||
if (checkCount % 6 == 0) {
|
if (checkCount % 6 == 0) {
|
||||||
log.info("第一阶段进度: 已接收={}, 已缓存={}, 设备数={}, 完成队列={}/4",
|
log.info("第一阶段进度: 已接收={}, 已缓存={}, 设备数={}, 完成队列={}/{}",
|
||||||
currentMessageCount, totalCached, messageCache.size(), completedQueues.size());
|
currentMessageCount, totalCached, messageCache.size(), completedQueues.size(), totalQueues);
|
||||||
}
|
}
|
||||||
|
|
||||||
// 检查是否有进度
|
// 检查是否有进度
|
||||||
@@ -342,7 +355,7 @@ public class TimeRangeReConsumer implements DisposableBean {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// 如果30秒内没有新消息,认为第一阶段完成
|
// 如果30秒内没有新消息,认为第一阶段完成
|
||||||
if (noProgressCount >= 6 && completedQueues.size() >= 4) {
|
if (noProgressCount >= 6 && completedQueues.size() >= totalQueues) {
|
||||||
log.info("30秒内无新消息且所有队列完成,第一阶段完成");
|
log.info("30秒内无新消息且所有队列完成,第一阶段完成");
|
||||||
shouldContinue.set(false);
|
shouldContinue.set(false);
|
||||||
break;
|
break;
|
||||||
@@ -542,65 +555,6 @@ public class TimeRangeReConsumer implements DisposableBean {
|
|||||||
log.info("✓ 重消费者状态已重置");
|
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<MessageQueue> 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();
|
totalCached += queue.size();
|
||||||
}
|
}
|
||||||
|
|
||||||
log.info("第一阶段进度: 已接收={}, 已缓存={}, 设备数={}, 完成队列={}/4",
|
log.info("第一阶段进度: 已接收={}, 已缓存={}, 设备数={}, 完成队列={}/{}",
|
||||||
receivedCount.get(), totalCached, messageCache.size(), completedQueues.size());
|
receivedCount.get(), totalCached, messageCache.size(), completedQueues.size(), totalQueues);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
Reference in New Issue
Block a user