跳到主要内容

基于Redis Zset的轻量级延迟任务调度器

· 阅读需 4 分钟

不知不觉又有好久没写了,于是还是抽空写篇。

最近项目里有个数据统计的需求,在用户提交数据后就更新下统计结果,但是,用户可能批量提交,为了性能考虑,需要稍微延迟下再进行统计,并且短时间内的重复提交只统计一次即可。

为了实现简便,最后基于本地缓存+Redis ZSet实现。

image

对于用户的提交,先使用本地缓存进行过滤,拦截瞬时大量提交;之后再经过Redis进行过滤,如果在一个时间窗口内提交多次,也只算一次。并且,由于应用的场景是提交后更新统计数据,所以要在时间窗口内最后一个提交之后再执行,而不是拦截掉后面的提交,这是和限流的区别。

还是先直接上代码吧。

public class StatRefreshService {

private static final String KEY_ZSET = "statRefresh:zset";
private static final String KEY_LOCK = "statRefresh:poll:lock";

@Resource
private RedissonClient redissonClient;

@Resource
@Qualifier("statRefreshExecutor")
private Executor statRefreshExecutor;


@SuppressWarnings("NullableProblems")
private Cache<String, String> localDedupCache;

private Thread pollThread;
private volatile boolean running = true;

@PostConstruct
public void init() {
this.localDedupCache = Caffeine.newBuilder()
.expireAfterWrite(120, TimeUnit.MILLISECONDS)
.maximumSize(500)
.build();

pollThread = new Thread(this::pollLoop, "stat-zset-poll");
pollThread.setDaemon(true);
pollThread.start();
log.info("StatRefreshService.init: 启动完成");
}

private void pollLoop() {
long pollInterval = 200;
while (running && !Thread.currentThread().isInterrupted()) {
try {
pullExpireTask();
} catch (Exception e) {
log.error("StatRefreshService 轮询拉取任务异常", e);
}
try {
//noinspection BusyWait
Thread.sleep(pollInterval);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
log.info("StatRefreshService ZSet轮询线程退出");
}

private void pullExpireTask() {
long nowTs = System.currentTimeMillis();
int batchSize = 20;
RScoredSortedSet<String> zset = redissonClient.getScoredSortedSet(KEY_ZSET);

Collection<ScoredEntry<String>> entryList = zset.entryRange(
Double.NEGATIVE_INFINITY, true,
(double) nowTs, true,
0, batchSize
);

if (entryList.isEmpty()) {
return;
}

RLock lock = redissonClient.getLock(KEY_LOCK);
boolean locked = false;
try {
locked = lock.tryLock(0, 1, TimeUnit.SECONDS);
} catch (InterruptedException e) {
log.error("StatRefreshService 获取锁异常", e);
Thread.currentThread().interrupt();
}
if (!locked) {
log.debug("StatRefreshService 未抢到轮询锁,本次跳过");
return;
}

try {
for (ScoredEntry<String> entry : entryList) {
String task = entry.getValue();
if (!zset.contains(task)) {
continue;
}
zset.remove(task);
log.info("StatRefreshService 拉取到期任务 task:{}", task);
statRefreshExecutor.execute(() -> {
try {
refreshStatData(task);
} catch (Exception e) {
log.error("StatRefreshService refreshStatData 异常, task:{}", task, e);
}
}
);
}
} catch (Exception e) {
log.error("StatRefreshService 处理到期任务异常", e);
} finally {
if (lock.isHeldByCurrentThread()) {
lock.unlock();
}
}
}

@PreDestroy
public void destroy() {
running = false;
if (pollThread != null) {
pollThread.interrupt();
}
log.info("StatRefreshService.destroy: 停止完成");
}

private void submit(String task, long delay) {
if (task == null || task.isBlank()) {
log.warn("StatRefreshService task 为空");
return;
}
String oldCache = localDedupCache.getIfPresent(task);
localDedupCache.invalidate(task);
localDedupCache.put(task, "1");
if (StringUtils.isNotEmpty(oldCache)) {
log.debug("StatRefreshService 本地拦截重复提交 task:{}", task);
return;
}

long now = System.currentTimeMillis();
long executeTime = now + delay;

RScoredSortedSet<String> zset = redissonClient.getScoredSortedSet(KEY_ZSET);
zset.add(executeTime, task);
log.debug("StatRefreshService 提交延迟任务 task:{} 执行时间戳:{}", task, executeTime);
}
}

提交任务的逻辑很直接,计算出执行时间戳作为score,把任务标识作为value塞进ZSet。

这里有个细节:ZSet的value是唯一的,如果同一个task重复提交,后面的会覆盖前面的score,相当于“刷新”了这个任务的执行时间。这个特性在我们的统计刷新场景里反而是好事——短时间内多次触发同一个统计项的刷新,只需要最后一次执行就够了。

通过轮询从Redis ZSet捞取到期任务,用entryRange按score范围查询,从负无穷到当前时间戳,分批捞取。

代码里还有一层Caffeine本地缓存,用来做提交去重,避免频繁操作Redis产生的不必要的网络IO。

@PostConstruct启动线程,@PreDestroy里标记停止位并中断线程。volatile关键字保证停止位的修改对轮询线程立即可见,不然可能出现主线程已经改了flag,但轮询线程迟迟看不到的情况。

当然,这个实现比较简单,还是存在一些问题,不过对于简单的场景足够用了。