From 24681c81c09022f584a57006f2534b5f74723414 Mon Sep 17 00:00:00 2001
From: 云 <2163098428@qq.com>
Date: 星期二, 30 六月 2026 09:27:31 +0800
Subject: [PATCH] 初始化项目

---
 yudao-framework/yudao-spring-boot-starter-mq/src/main/java/cn/iocoder/yudao/framework/mq/redis/core/job/RedisStreamMessageCleanupJob.java |   30 +++++++++++++++++++++++-------
 1 files changed, 23 insertions(+), 7 deletions(-)

diff --git a/yudao-framework/yudao-spring-boot-starter-mq/src/main/java/cn/iocoder/yudao/framework/mq/redis/core/job/RedisStreamMessageCleanupJob.java b/yudao-framework/yudao-spring-boot-starter-mq/src/main/java/cn/iocoder/yudao/framework/mq/redis/core/job/RedisStreamMessageCleanupJob.java
index 19da845..eff64fd 100644
--- a/yudao-framework/yudao-spring-boot-starter-mq/src/main/java/cn/iocoder/yudao/framework/mq/redis/core/job/RedisStreamMessageCleanupJob.java
+++ b/yudao-framework/yudao-spring-boot-starter-mq/src/main/java/cn/iocoder/yudao/framework/mq/redis/core/job/RedisStreamMessageCleanupJob.java
@@ -23,7 +23,16 @@
 @AllArgsConstructor
 public class RedisStreamMessageCleanupJob {
 
-    private static final String LOCK_KEY = "redis:stream:message-cleanup:lock";
+    /**
+     * 涓氬姟 MQ锛圫pring 瀹瑰櫒鍐� AbstractRedisStreamMessageListener锛夋竻鐞嗕换鍔′娇鐢ㄧ殑鍒嗗竷寮忛攣
+     */
+    public static final String DEFAULT_CLEANUP_LOCK_KEY = "redis:stream:message-cleanup:lock";
+
+    /**
+     * IoT Redis 鎬荤嚎娓呯悊浠诲姟浣跨敤鐨勫垎甯冨紡閿侊紙椤讳笌 {@link #DEFAULT_CLEANUP_LOCK_KEY} 鍖哄垎锛屽惁鍒欎細鍏辨姠涓�鎶婇攣锛�
+     * 鍚屼竴鏃跺埢鍙湁涓�渚ц兘鎵ц XTRIM锛屽彟涓�渚� Stream 鍙兘鏃犻檺绉帇锛�
+     */
+    public static final String IOT_CLEANUP_LOCK_KEY = "redis:stream:message-cleanup:lock:iot";
 
     /**
      * 淇濈暀鐨勬秷鎭暟閲忥紝榛樿淇濈暀鏈�杩� 10000 鏉℃秷鎭�
@@ -33,22 +42,29 @@
     private final List<AbstractRedisStreamMessageListener<?>> listeners;
     private final RedisMQTemplate redisTemplate;
     private final RedissonClient redissonClient;
+    /**
+     * Redisson 閿侀敭锛堝 Bean 娉ㄥ唽娓呯悊浠诲姟鏃跺繀椤诲悇涓嶇浉鍚岋級
+     */
+    private final String cleanupLockKey;
 
     /**
      * 姣忓皬鏃舵墽琛屼竴娆℃竻鐞嗕换鍔�
      */
     @Scheduled(cron = "0 0 * * * ?")
     public void cleanup() {
-        RLock lock = redissonClient.getLock(LOCK_KEY);
-        // 灏濊瘯鍔犻攣
+        RLock lock = redissonClient.getLock(cleanupLockKey);
         if (lock.tryLock()) {
             try {
                 execute();
             } catch (Exception ex) {
-                log.error("[cleanup][鎵ц寮傚父]", ex);
+                log.error("[cleanup][鎵ц寮傚父][lockKey={}]", cleanupLockKey, ex);
             } finally {
-                lock.unlock();
+                if (lock.isHeldByCurrentThread()) {
+                    lock.unlock();
+                }
             }
+        } else {
+            log.debug("[cleanup][鏈幏鍙栧埌閿侊紝璺宠繃鏈疆][lockKey={}]", cleanupLockKey);
         }
     }
 
@@ -59,8 +75,8 @@
         StreamOperations<String, Object, Object> ops = redisTemplate.getRedisTemplate().opsForStream();
         listeners.forEach(listener -> {
             try {
-                // 浣跨敤 XTRIM 鍛戒护娓呯悊娑堟伅锛屽彧淇濈暀鏈�杩戠殑 MAX_LEN 鏉℃秷鎭�
-                Long trimCount = ops.trim(listener.getStreamKey(), MAX_COUNT, true);
+                // 浣跨敤 XTRIM MAXLEN 绮剧‘瑁佸壀锛坅pproximate=false锛夛紝閬垮厤 ~ 妯″紡涓嬮暱鏈熸槑鏄鹃珮浜庝笂闄�
+                Long trimCount = ops.trim(listener.getStreamKey(), MAX_COUNT, false);
                 if (trimCount != null && trimCount > 0) {
                     log.info("[execute][Stream({}) 娓呯悊娑堟伅鏁伴噺({})]", listener.getStreamKey(), trimCount);
                 }

--
Gitblit v1.9.3