From 67a627054bb4f9ebaa6833994e189ec643e2b803 Mon Sep 17 00:00:00 2001 From: gaoyutao Date: Sat, 28 Feb 2026 18:01:30 +0800 Subject: [PATCH] =?UTF-8?q?=E6=B5=81=E9=87=8F=E8=A1=A5=E5=85=A8=E6=B6=88?= =?UTF-8?q?=E8=B4=B9=E8=80=85=E6=94=B9=E4=B8=BA=E7=BC=93=E5=AD=98=E6=8E=92?= =?UTF-8?q?=E5=BA=8F=E5=A4=84=E7=90=86=EF=BC=8C=E4=BF=9D=E8=AF=81=E9=A1=BA?= =?UTF-8?q?=E5=BA=8F=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../rocketmq/handler/MessageHandler.java | 81 ++++++++++++------- 1 file changed, 52 insertions(+), 29 deletions(-) diff --git a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/handler/MessageHandler.java b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/handler/MessageHandler.java index e76e1a1..e306d21 100644 --- a/tongran-rocketmq/src/main/java/com/tongran/rocketmq/handler/MessageHandler.java +++ b/tongran-rocketmq/src/main/java/com/tongran/rocketmq/handler/MessageHandler.java @@ -28,7 +28,6 @@ import org.apache.http.impl.client.CloseableHttpClient; import org.apache.http.impl.client.HttpClients; import org.apache.http.impl.conn.PoolingHttpClientConnectionManager; import org.apache.http.util.EntityUtils; -import org.redisson.api.RLock; import org.redisson.api.RedissonClient; import org.springframework.beans.BeanUtils; import org.springframework.beans.factory.annotation.Autowired; @@ -90,7 +89,8 @@ public class MessageHandler { private static final String USED_PORTS_KEY = "frpc:used_ports"; private static final long HEARTBEAT_TIMEOUT = 30000; // 3分钟超时 - + private final ConcurrentHashMap> messageBuffer = + new ConcurrentHashMap<>(); @Autowired private RedisTemplate redisTemplate; @Autowired @@ -875,36 +875,59 @@ public class MessageHandler { } private void handleNetRecoverMessage(DeviceMessage message) { String clientId = message.getClientId(); - String lockKey = "traffic:recover:" + clientId; - RLock lock = redissonClient.getLock(lockKey); - boolean locked = false; + List interfaces = JsonDataParser.parseJsonData(message.getData(), InitialBandwidthTraffic.class); - try { - // 尝试获取锁 - locked = lock.tryLock(0, 20, TimeUnit.SECONDS); - if (locked) { -// log.info("设备{}获取锁成功,开始处理消息", clientId); - processNetRecoverMessageInternal(message); -// log.info("设备{}消息处理完成", clientId); - } else { - log.warn("设备{}获取锁失败,消息处理繁忙", clientId); - throw new RuntimeException("设备处理繁忙,请重试"); - } - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - log.error("处理消息时线程被中断", e); - } finally { - // 只有成功获取锁的线程才需要释放 - if (locked && lock.isHeldByCurrentThread()) { - try { - lock.unlock(); -// log.debug("设备{}锁已释放", clientId); - } catch (IllegalMonitorStateException e) { - log.warn("设备{}锁释放异常,可能已自动超时", clientId); - } - } + if (interfaces.isEmpty()) { + return; + } + + boolean isLast = interfaces.get(0).isLastTrafficFlag(); + + // 1. 消息入队 + LinkedBlockingQueue queue = messageBuffer + .computeIfAbsent(clientId, k -> new LinkedBlockingQueue<>()); + queue.offer(message); + + // 2. 如果是最后一条消息,开始处理(此时前面的消息肯定都已经到了) + if (isLast) { + processAllMessagesInOrder(clientId); } } + + private void processAllMessagesInOrder(String clientId) { + LinkedBlockingQueue queue = messageBuffer.get(clientId); + if (queue == null || queue.isEmpty()) { + return; + } + + // 1. 取出所有消息 + List allMessages = new ArrayList<>(); + queue.drainTo(allMessages); + + // 2. 按时间戳排序 + allMessages.sort((msg1, msg2) -> { + long ts1 = JsonDataParser.parseJsonData(msg1.getData(), InitialBandwidthTraffic.class) + .get(0).getTimestamp(); + long ts2 = JsonDataParser.parseJsonData(msg2.getData(), InitialBandwidthTraffic.class) + .get(0).getTimestamp(); + return Long.compare(ts1, ts2); + }); + + // 3. 顺序处理 + for (DeviceMessage message : allMessages) { + try { + processNetRecoverMessageInternal(message); + } catch (Exception e) { + log.error("处理设备{}消息失败", clientId, e); + } + } + + log.info("设备{}顺序处理完成,共处理{}条消息", clientId, allMessages.size()); + + // 4. 清理 + messageBuffer.remove(clientId); + } + /** * 网络重试流量数据入库 * @param message