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);
}