Browse Source

fix(redis): 生产路径禁止使用 KEYS,改用 SCAN,修复 Redis 阻塞导致巡检服务不可用

现象:巡检服务(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,对删除/遍历场景无影响),未改动任何业务判定逻辑。
yanyuan
wangguangyuan 2 weeks ago
parent
commit
de77ada07a
7 changed files with 96 additions and 27 deletions
  1. +51
    -0
      inspect-base/inspect-base-redis/src/main/java/com/inspect/base/redis/service/RedisService.java
  2. +2
    -2
      inspect-job/src/main/java/com/inspect/job/task/JobMainTask.java
  3. +35
    -19
      inspect-main/inspect-main-task-exec/src/main/java/com/inspect/exec/controller/PatrolTaskExecController.java
  4. +5
    -3
      inspect-main/inspect-main-task/src/main/java/com/inspect/task/controller/PatrolTaskController.java
  5. +1
    -1
      inspect-main/inspect-main-video/src/main/java/com/inspect/patrol/controller/TestController.java
  6. +1
    -1
      inspect-management/src/main/java/com/inspect/system/controller/SysUserOnlineController.java
  7. +1
    -1
      inspect-management/src/main/java/com/inspect/system/service/impl/SysConfigServiceImpl.java

+ 51
- 0
inspect-base/inspect-base-redis/src/main/java/com/inspect/base/redis/service/RedisService.java View File

@ -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非阻塞
*
* <p>禁止在业务代码里使用 {@link #keys(String)}KEYS O(N) 的阻塞命令 Redis 是单线程的
* 一次 KEYS 执行期间其它所有客户端的命令都要排队keyspace 较大时单次 KEYS 需要十几毫秒
* 若在高频路径上循环调用会直接把 Redis 及其上游服务拖死已在生产环境造成过故障
* 需要按 pattern key 时一律使用本方法
*
* @param pattern 匹配模式例如 TASK_CODE@279@*
* @return 匹配到的 key 集合SCAN 可能重复返回同一个 key方法内部已去重
*/
public Collection<String> scan(String pattern) {
return scan(pattern, 500);
}
/**
* @param count 每批返回的 key 数量提示值COUNT非硬性上限
*/
public Collection<String> scan(String pattern, int count) {
Set<String> 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<Void>) connection -> {
try (Cursor<byte[]> 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<String> keys(String pattern) {
return redisTemplate.keys(pattern);
}


+ 2
- 2
inspect-job/src/main/java/com/inspect/job/task/JobMainTask.java View File

@ -2074,7 +2074,7 @@ public class JobMainTask {
Date date = new Date();
log.info("JobMainTask fixedTimer clear taskInfos: {}", taskInfos.size());
taskInfos.clear();
Collection<String> redisKeys = redisService.keys(RedisConst.TASK_CODE_EX);
Collection<String> 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<String> redisKeys = redisService.keys(key);
Collection<String> redisKeys = redisService.scan(key);
log.info("TodayTask key:{}, size:{}", key, redisKeys.size());
for (String redisKey : redisKeys) {
try {


+ 35
- 19
inspect-main/inspect-main-task-exec/src/main/java/com/inspect/exec/controller/PatrolTaskExecController.java View File

@ -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<Date> exeTimes = getExecTimes(task.getIntervalExecuteTime(), DateUtils.format(DateUtils.yyyyMMddHHmmss2, task.getIntervalStartTime()), DateUtils.format(DateUtils.yyyyMMddHHmmss2, task.getIntervalEndTime()), intervalNumber);
// List<Date> 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<Date> 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<String> devNoList = new ArrayList<>();
@ -646,18 +649,6 @@ public class PatrolTaskExecController extends BaseController {
if (taskExecRecord == null) {
List<PatrolTask> patrolTasks = getPatrolTasks(task);
if (!patrolTasks.isEmpty()) {
if(finalDate != null) {
Collection<String> 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
*
* <p>历史实现放在 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<String> 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<Date> 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());


+ 5
- 3
inspect-main/inspect-main-task/src/main/java/com/inspect/task/controller/PatrolTaskController.java View File

@ -586,8 +586,10 @@ public class PatrolTaskController extends BaseController {
patrolTaskService.deletePatrolTaskByTaskIds(taskIds);
patrolTaskInfoService.deletePatrolTaskInfoByMajorIds(taskIds);
// 原来每个 taskId 都做一次 KEYS 全库扫描改为只扫描一次SCAN 非阻塞
Collection<String> 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<String> redisKeys = redisService.keys(RedisConst.TASK_CODE + patrolTask.getTaskCode() + StringUtils.AT + "*");
Collection<String> 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<String> redisKeys = redisService.keys(RedisConst.TASK_CODE + patrolTask.getTaskCode() + StringUtils.AT + "*");
Collection<String> redisKeys = redisService.scan(RedisConst.TASK_CODE + patrolTask.getTaskCode() + StringUtils.AT + "*");
for (String redisKey : redisKeys) {
redisService.deleteObject(redisKey);
}


+ 1
- 1
inspect-main/inspect-main-video/src/main/java/com/inspect/patrol/controller/TestController.java View File

@ -98,7 +98,7 @@ public class TestController extends BaseController {
this.redisService.deleteObject(REDIS_LINKAGE_QUEUE_NOW);
for (int i = 0; i < 10; ++i) {
Collection<String> keys = this.redisService.keys(i + "*");
Collection<String> keys = this.redisService.scan(i + "*");
this.redisService.deleteObject(keys);
total += keys.size();
}


+ 1
- 1
inspect-management/src/main/java/com/inspect/system/controller/SysUserOnlineController.java View File

@ -39,7 +39,7 @@ public class SysUserOnlineController extends BaseController {
@RequiresPermissions({"monitor:online:list"})
@GetMapping({"/list"})
public TableDataInfo list(String ipaddr, String userName) {
Collection<String> keys = this.redisService.keys("login_tokens:*");
Collection<String> keys = this.redisService.scan("login_tokens:*");
List<SysUserOnline> userOnlineList = new ArrayList<>();
for (String key : keys) {
LoginUser user = this.redisService.getCacheObject(key);


+ 1
- 1
inspect-management/src/main/java/com/inspect/system/service/impl/SysConfigServiceImpl.java View File

@ -103,7 +103,7 @@ public class SysConfigServiceImpl implements ISysConfigService {
}
public void clearConfigCache() {
Collection<String> keys = this.redisService.keys("sys_config:*");
Collection<String> keys = this.redisService.scan("sys_config:*");
this.redisService.deleteObject(keys);
}


Loading…
Cancel
Save