From de77ada07a481bee19f5bb00354742fa6d58c05c Mon Sep 17 00:00:00 2001 From: wangguangyuan Date: Sat, 12 Sep 2026 14:19:48 +0800 Subject: [PATCH] =?UTF-8?q?fix(redis):=20=E7=94=9F=E4=BA=A7=E8=B7=AF?= =?UTF-8?q?=E5=BE=84=E7=A6=81=E6=AD=A2=E4=BD=BF=E7=94=A8=20KEYS=EF=BC=8C?= =?UTF-8?q?=E6=94=B9=E7=94=A8=20SCAN=EF=BC=8C=E4=BF=AE=E5=A4=8D=20Redis=20?= =?UTF-8?q?=E9=98=BB=E5=A1=9E=E5=AF=BC=E8=87=B4=E5=B7=A1=E6=A3=80=E6=9C=8D?= =?UTF-8?q?=E5=8A=A1=E4=B8=8D=E5=8F=AF=E7=94=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 现象:巡检服务(dliip-patrol)容器连续约 9 小时 unhealthy;Tomcat 200 个请求线程 全部阻塞在 Lettuce(AsyncCommand.await 无超时,永不失败),容器 CPU 在 2%~500% 间跳变, 健康检查(curl 127.0.0.1:9901,10s 超时)连续 284 次超时。 根因:inspect-job 定时调用 /exec/makeCurrentDayTask,该接口遍历所有启用任务,对每个 "间隔执行"任务的每个执行时刻调用一次 parseTaskToRedis,而 parseTaskToRedis 内部用 redisService.keys(TASK_CODE@taskCode@*) 做全库扫描来清理策略 key。于是一次调用会对 Redis 发起成百上千次 KEYS:Redis 单线程,单次 KEYS 约 13ms 且返回体巨大(实测输出 8MB/s,单个客户端连接积压 321 条回复 / 5.8MB 输出缓冲),KEYS 执行期间其它所有客户端 的命令全部排队,实测 Redis 响应延迟 avg 287ms / max 1365ms,进而导致上游业务请求堆积、 Tomcat 线程池被占满、健康检查排不上队而超时。 修复: 1. RedisService 新增 scan(pattern[, count]),基于 SCAN 游标分批读取,不阻塞其它客户端; 原 keys() 标记 @Deprecated 并注明生产环境禁用。 2. 全仓库 11 处 redisService.keys() 调用点改为 scan()。 3. PatrolTaskExecController:抽出 cleanTaskKeysAfter(),策略 key 清理由"每个执行时刻一次" 改为"每个任务一次",一次 makeCurrentDayTask 的 KEYS 次数从成百上千次降到与任务数同量级。 4. PatrolTaskController:批量删除任务时原为 N 个 taskId 扫描 N 次全库,改为只扫描一次。 影响面:仅改变"按 pattern 查找 key"的实现方式(KEYS→SCAN,语义等价;SCAN 可能重复返回 同一个 key,对删除/遍历场景无影响),未改动任何业务判定逻辑。 --- .../base/redis/service/RedisService.java | 51 ++++++++++++++++++ .../com/inspect/job/task/JobMainTask.java | 4 +- .../controller/PatrolTaskExecController.java | 54 ++++++++++++------- .../task/controller/PatrolTaskController.java | 8 +-- .../patrol/controller/TestController.java | 2 +- .../controller/SysUserOnlineController.java | 2 +- .../service/impl/SysConfigServiceImpl.java | 2 +- 7 files changed, 96 insertions(+), 27 deletions(-) diff --git a/inspect-base/inspect-base-redis/src/main/java/com/inspect/base/redis/service/RedisService.java b/inspect-base/inspect-base-redis/src/main/java/com/inspect/base/redis/service/RedisService.java index faa8294..d1f81d1 100644 --- a/inspect-base/inspect-base-redis/src/main/java/com/inspect/base/redis/service/RedisService.java +++ b/inspect-base/inspect-base-redis/src/main/java/com/inspect/base/redis/service/RedisService.java @@ -2,6 +2,7 @@ package com.inspect.base.redis.service; import java.util.Collection; import java.util.Iterator; +import java.util.LinkedHashSet; import java.util.List; import java.util.Map; import java.util.Set; @@ -11,8 +12,11 @@ import com.inspect.base.core.constant.Color; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.data.redis.core.BoundSetOperations; +import org.springframework.data.redis.core.Cursor; import org.springframework.data.redis.core.HashOperations; +import org.springframework.data.redis.core.RedisCallback; import org.springframework.data.redis.core.RedisTemplate; +import org.springframework.data.redis.core.ScanOptions; import org.springframework.data.redis.core.ValueOperations; import org.springframework.stereotype.Component; @@ -135,6 +139,53 @@ public class RedisService { return redisTemplate.opsForHash().multiGet(key, hKeys); } + /** + * 游标分批扫描 key(非阻塞)。 + * + *

禁止在业务代码里使用 {@link #keys(String)}:KEYS 是 O(N) 的阻塞命令,而 Redis 是单线程的, + * 一次 KEYS 执行期间其它所有客户端的命令都要排队。keyspace 较大时单次 KEYS 需要十几毫秒, + * 若在高频路径上循环调用,会直接把 Redis 及其上游服务拖死(已在生产环境造成过故障)。 + * 需要按 pattern 查 key 时一律使用本方法。 + * + * @param pattern 匹配模式,例如 TASK_CODE@279@* + * @return 匹配到的 key 集合(SCAN 可能重复返回同一个 key,方法内部已去重) + */ + public Collection scan(String pattern) { + return scan(pattern, 500); + } + + /** + * @param count 每批返回的 key 数量提示值(COUNT,非硬性上限) + */ + public Collection scan(String pattern, int count) { + Set keys = new LinkedHashSet<>(); + if (pattern == null || pattern.isEmpty()) { + return keys; + } + + long start = System.currentTimeMillis(); + ScanOptions options = ScanOptions.scanOptions().match(pattern).count(count).build(); + redisTemplate.execute((RedisCallback) connection -> { + try (Cursor cursor = connection.scan(options)) { + while (cursor.hasNext()) { + Object key = redisTemplate.getKeySerializer().deserialize(cursor.next()); + if (key != null) { + keys.add(String.valueOf(key)); + } + } + } catch (Exception e) { + logger.error("RedisService.scan error, pattern: {}", pattern, e); + } + return null; + }); + logger.info("RedisService.scan pattern: {}, size: {}, cost: {} ms", pattern, keys.size(), System.currentTimeMillis() - start); + return keys; + } + + /** + * @deprecated KEYS 会阻塞整个 Redis(单线程),生产环境禁止调用,请改用 {@link #scan(String)}。 + */ + @Deprecated public Collection keys(String pattern) { return redisTemplate.keys(pattern); } diff --git a/inspect-job/src/main/java/com/inspect/job/task/JobMainTask.java b/inspect-job/src/main/java/com/inspect/job/task/JobMainTask.java index 30d3198..b5fd71a 100644 --- a/inspect-job/src/main/java/com/inspect/job/task/JobMainTask.java +++ b/inspect-job/src/main/java/com/inspect/job/task/JobMainTask.java @@ -2074,7 +2074,7 @@ public class JobMainTask { Date date = new Date(); log.info("JobMainTask fixedTimer clear taskInfos: {}", taskInfos.size()); taskInfos.clear(); - Collection redisKeys = redisService.keys(RedisConst.TASK_CODE_EX); + Collection redisKeys = redisService.scan(RedisConst.TASK_CODE_EX); for (String redisKey : redisKeys) { String[] keywords = StringUtils.split(redisKey, StringUtils.AT); if (keywords.length == 3) { @@ -2093,7 +2093,7 @@ public class JobMainTask { log.debug("***************************** JobTaskTimer execEveryDayTask *************************************"); String today = LocalDate.now().toString(); String key = RedisConst.TASK_CODE_EX + today + "*"; - Collection redisKeys = redisService.keys(key); + Collection redisKeys = redisService.scan(key); log.info("TodayTask key:{}, size:{}", key, redisKeys.size()); for (String redisKey : redisKeys) { try { diff --git a/inspect-main/inspect-main-task-exec/src/main/java/com/inspect/exec/controller/PatrolTaskExecController.java b/inspect-main/inspect-main-task-exec/src/main/java/com/inspect/exec/controller/PatrolTaskExecController.java index e377c16..6387a3d 100644 --- a/inspect-main/inspect-main-task-exec/src/main/java/com/inspect/exec/controller/PatrolTaskExecController.java +++ b/inspect-main/inspect-main-task-exec/src/main/java/com/inspect/exec/controller/PatrolTaskExecController.java @@ -546,7 +546,7 @@ public class PatrolTaskExecController extends BaseController { task.setFixedStartTime(DateUtils.parse(DateUtils.yyyyMMddHHmmss2, (DateUtils.format(DateUtils.yyyyMMdd2, new Date()) + " " + cycleTimes[i]))); final String taskType = "CYCLE-BY-WEEK"; - parseTaskToRedis(taskType, task, null); + parseTaskToRedis(taskType, task); } } else if (isCycleTaskByMonth(task)) { String[] monthList = task.getCycleMonth().split(StringUtils.COMMA); @@ -559,7 +559,7 @@ public class PatrolTaskExecController extends BaseController { task.setFixedStartTime(DateUtils.parse(DateUtils.yyyyMMddHHmmss2, (DateUtils.format(DateUtils.yyyyMMdd2, new Date()) + " " + cycleTimes[i]))); final String taskType = "CYCLE-BY-MONTH"; - parseTaskToRedis(taskType, task, null); + parseTaskToRedis(taskType, task); } } } else if (isInterTask(task)) { @@ -571,11 +571,13 @@ public class PatrolTaskExecController extends BaseController { if (intervalNumber > 0) { List exeTimes = getExecTimes(task.getIntervalExecuteTime(), DateUtils.format(DateUtils.yyyyMMddHHmmss2, task.getIntervalStartTime()), DateUtils.format(DateUtils.yyyyMMddHHmmss2, task.getIntervalEndTime()), intervalNumber); // List exeTimes = getExecTimesOld(task.getIntervalExecuteTime(), intervalNumber); + // 每个任务的策略 key 只清理一次(原实现在每个执行时刻都做一次 KEYS 全库扫描) + cleanTaskKeysAfter(task, task.getIntervalEndTime()); for (Date exeTime : exeTimes) { logger.debug("[TASK] {}, isInterTaskByHour exeTime: {}", task.getTaskCode(), DateUtils.format(DateUtils.yyyyMMddHHmmss2, exeTime)); task.setFixedStartTime(exeTime); final String taskType = "INTER-BY-HOUR"; - parseTaskToRedis(taskType, task, task.getIntervalEndTime()); + parseTaskToRedis(taskType, task); } } else { log.info("[TASK] isInterTaskByHour intervalNumber error: {}", intervalNumber); @@ -589,10 +591,11 @@ public class PatrolTaskExecController extends BaseController { List exeTimes = getExecTimesByMinute(task.getIntervalExecuteTime(), DateUtils.format(DateUtils.yyyyMMddHHmmss2, task.getIntervalStartTime()), endTime, intervalNumber); if (exeTimes.size() > 0) { // Date date = exeTimes.get(exeTimes.size() - 1); + cleanTaskKeysAfter(task, task.getIntervalEndTime()); for (Date exeTime : exeTimes) { task.setFixedStartTime(exeTime); final String taskType = "INTER-BY-MINUTE"; - parseTaskToRedis(taskType, task, task.getIntervalEndTime()); + parseTaskToRedis(taskType, task); } } } else { @@ -607,7 +610,7 @@ public class PatrolTaskExecController extends BaseController { String time = intervalExecuteTime[0] + ":" + DateUtils.getMinuteInt() + ":" + intervalExecuteTime[2]; task.setFixedStartTime(DateUtils.parse(DateUtils.yyyyMMddHHmmss2, (DateUtils.format(DateUtils.yyyyMMdd2, new Date()) + " " + time))); final String taskType = "INTER-BY-DATE"; - parseTaskToRedis(taskType, task, null); + parseTaskToRedis(taskType, task); } } } @@ -619,7 +622,7 @@ public class PatrolTaskExecController extends BaseController { } } - private void parseTaskToRedis(final String taskType, PatrolTask task, Date finalDate) { + private void parseTaskToRedis(final String taskType, PatrolTask task) { if (StringUtils.isNotEmpty(task.getDevNo())) { String[] devNos = task.getDevNo().split(StringUtils.COMMA); List devNoList = new ArrayList<>(); @@ -646,18 +649,6 @@ public class PatrolTaskExecController extends BaseController { if (taskExecRecord == null) { List patrolTasks = getPatrolTasks(task); if (!patrolTasks.isEmpty()) { - if(finalDate != null) { - Collection redisKeys = redisService.keys(RedisConst.TASK_CODE + task.getTaskCode() + StringUtils.AT + "*"); - for (String redisKey : redisKeys) { - String[] keywords = StringUtils.split(redisKey, StringUtils.AT); - if (keywords.length == 3) { - String fixedStartTime = keywords[2]; - if (DateUtils.parse(DateUtils.yyyyMMddHHmmss2, fixedStartTime).after(finalDate)) { - redisService.deleteObject(redisKey); - } - } - } - } logger.debug(Color.GREEN + "[TASK] TYPE: {}, CYCLE BY WEEK key: {}, patrolId: {}" + Color.END, taskType, key, taskPatrolledId); redisService.setCacheObject(key, JSONArray.toJSONString(patrolTasks)); } @@ -667,6 +658,31 @@ public class PatrolTaskExecController extends BaseController { } } + /** + * 清理该任务下"计划执行时间晚于 finalDate"的策略 key。 + * + *

历史实现放在 parseTaskToRedis 里,而 parseTaskToRedis 是按"每个执行时刻"调用的, + * 于是一次 makeCurrentDayTask 会对 Redis 发起成百上千次 KEYS 全库扫描(TASK_CODE@taskCode@*), + * 单线程 Redis 被阻塞,导致所有业务请求排队、tomcat 线程池被占满、容器健康检查连续超时。 + * 现在改为:① 用 SCAN 替代 KEYS(不阻塞其它客户端);② 每个任务只清理一次。 + */ + private void cleanTaskKeysAfter(PatrolTask task, Date finalDate) { + if (task == null || finalDate == null || StringUtils.isEmpty(task.getTaskCode())) { + return; + } + + Collection redisKeys = redisService.scan(RedisConst.TASK_CODE + task.getTaskCode() + StringUtils.AT + "*"); + for (String redisKey : redisKeys) { + String[] keywords = StringUtils.split(redisKey, StringUtils.AT); + if (keywords.length == 3) { + Date fixedStartTime = DateUtils.parse(DateUtils.yyyyMMddHHmmss2, keywords[2]); + if (fixedStartTime != null && fixedStartTime.after(finalDate)) { + redisService.deleteObject(redisKey); + } + } + } + } + private List getExecTimesOld(String exeTime, int num) { SimpleDateFormat sdf1 = new SimpleDateFormat("yyyy-MM-dd"); SimpleDateFormat sdfTime = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss"); @@ -795,7 +811,7 @@ public class PatrolTaskExecController extends BaseController { // this.log.info(" 【当前正在执行的任务:{}, 其它任务不能执行!!!】", taskCodeStr); // } else { SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss"); - for (String key : this.redisService.keys(RedisConst.TASK_CODE_EX)) { + for (String key : this.redisService.scan(RedisConst.TASK_CODE_EX)) { String fixedStartTime = key.split(StringUtils.AT)[2]; long currentMinutes = TimeUnit.MILLISECONDS.toMinutes(System.currentTimeMillis()); long fixedStartMinutes = TimeUnit.MILLISECONDS.toMinutes(sdf.parse(fixedStartTime).getTime()); diff --git a/inspect-main/inspect-main-task/src/main/java/com/inspect/task/controller/PatrolTaskController.java b/inspect-main/inspect-main-task/src/main/java/com/inspect/task/controller/PatrolTaskController.java index 2c4306c..de1fbba 100644 --- a/inspect-main/inspect-main-task/src/main/java/com/inspect/task/controller/PatrolTaskController.java +++ b/inspect-main/inspect-main-task/src/main/java/com/inspect/task/controller/PatrolTaskController.java @@ -586,8 +586,10 @@ public class PatrolTaskController extends BaseController { patrolTaskService.deletePatrolTaskByTaskIds(taskIds); patrolTaskInfoService.deletePatrolTaskInfoByMajorIds(taskIds); + // 原来每个 taskId 都做一次 KEYS 全库扫描,改为只扫描一次(SCAN 非阻塞) + Collection taskCodeRedisKeys = redisService.scan(RedisConst.TASK_CODE_EX); for (Long taskId : taskIds) { - for (String key : redisService.keys(RedisConst.TASK_CODE_EX)) { + for (String key : taskCodeRedisKeys) { logger.info("[CONTROLLER] remove task redis key: {}", key); if (key.contains(String.valueOf(taskId))) { redisService.deleteObject(key); @@ -1118,7 +1120,7 @@ public class PatrolTaskController extends BaseController { issueTask(patrolTask); // 修改任务,删除原有的任务执行策略,根据新的策略重新生成 if (StringUtils.isNotEmpty(patrolTask.getTaskCode())) { - Collection redisKeys = redisService.keys(RedisConst.TASK_CODE + patrolTask.getTaskCode() + StringUtils.AT + "*"); + Collection redisKeys = redisService.scan(RedisConst.TASK_CODE + patrolTask.getTaskCode() + StringUtils.AT + "*"); for (String redisKey : redisKeys) { redisService.deleteObject(redisKey); } @@ -1344,7 +1346,7 @@ public class PatrolTaskController extends BaseController { issueTask(patrolTask); // 修改任务,删除原有的任务执行策略,根据新的策略重新生成 if (StringUtils.isNotEmpty(patrolTask.getTaskCode())) { - Collection redisKeys = redisService.keys(RedisConst.TASK_CODE + patrolTask.getTaskCode() + StringUtils.AT + "*"); + Collection redisKeys = redisService.scan(RedisConst.TASK_CODE + patrolTask.getTaskCode() + StringUtils.AT + "*"); for (String redisKey : redisKeys) { redisService.deleteObject(redisKey); } diff --git a/inspect-main/inspect-main-video/src/main/java/com/inspect/patrol/controller/TestController.java b/inspect-main/inspect-main-video/src/main/java/com/inspect/patrol/controller/TestController.java index ae6cd64..6b3cb17 100644 --- a/inspect-main/inspect-main-video/src/main/java/com/inspect/patrol/controller/TestController.java +++ b/inspect-main/inspect-main-video/src/main/java/com/inspect/patrol/controller/TestController.java @@ -98,7 +98,7 @@ public class TestController extends BaseController { this.redisService.deleteObject(REDIS_LINKAGE_QUEUE_NOW); for (int i = 0; i < 10; ++i) { - Collection keys = this.redisService.keys(i + "*"); + Collection keys = this.redisService.scan(i + "*"); this.redisService.deleteObject(keys); total += keys.size(); } diff --git a/inspect-management/src/main/java/com/inspect/system/controller/SysUserOnlineController.java b/inspect-management/src/main/java/com/inspect/system/controller/SysUserOnlineController.java index db0359f..bbcec29 100644 --- a/inspect-management/src/main/java/com/inspect/system/controller/SysUserOnlineController.java +++ b/inspect-management/src/main/java/com/inspect/system/controller/SysUserOnlineController.java @@ -39,7 +39,7 @@ public class SysUserOnlineController extends BaseController { @RequiresPermissions({"monitor:online:list"}) @GetMapping({"/list"}) public TableDataInfo list(String ipaddr, String userName) { - Collection keys = this.redisService.keys("login_tokens:*"); + Collection keys = this.redisService.scan("login_tokens:*"); List userOnlineList = new ArrayList<>(); for (String key : keys) { LoginUser user = this.redisService.getCacheObject(key); diff --git a/inspect-management/src/main/java/com/inspect/system/service/impl/SysConfigServiceImpl.java b/inspect-management/src/main/java/com/inspect/system/service/impl/SysConfigServiceImpl.java index c90469f..a407d1f 100644 --- a/inspect-management/src/main/java/com/inspect/system/service/impl/SysConfigServiceImpl.java +++ b/inspect-management/src/main/java/com/inspect/system/service/impl/SysConfigServiceImpl.java @@ -103,7 +103,7 @@ public class SysConfigServiceImpl implements ISysConfigService { } public void clearConfigCache() { - Collection keys = this.redisService.keys("sys_config:*"); + Collection keys = this.redisService.scan("sys_config:*"); this.redisService.deleteObject(keys); }