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