update 增强 SSE 连接管理,支持延迟续期和心跳检测配置

This commit is contained in:
AprilWind 2025-10-17 21:16:54 +08:00
parent 0849c2e1c1
commit 57f4eaa952
3 changed files with 73 additions and 17 deletions

View File

@ -3,6 +3,8 @@ package org.dromara.common.sse.config;
import lombok.Data; import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.boot.context.properties.ConfigurationProperties;
import java.util.concurrent.TimeUnit;
/** /**
* SSE 配置项 * SSE 配置项
* *
@ -12,10 +14,29 @@ import org.springframework.boot.context.properties.ConfigurationProperties;
@ConfigurationProperties("sse") @ConfigurationProperties("sse")
public class SseProperties { public class SseProperties {
/**
* 是否启用 SSE 功能
*/
private Boolean enabled; private Boolean enabled;
/** /**
* 路径 * 路径
*/ */
private String path; private String path;
/**
* 心跳检测间隔默认秒
*/
private long heartbeatSeconds = 30;
/**
* 清理线程执行间隔默认秒
*/
private long checkIntervalSeconds = 10;
/**
* 时间单位分钟等默认秒
*/
private TimeUnit unit = TimeUnit.SECONDS;
} }

View File

@ -41,7 +41,7 @@ public class SseEmitterDelayed implements Delayed {
/** /**
* 该连接的过期时间戳毫秒到达该时间后可视为超时 * 该连接的过期时间戳毫秒到达该时间后可视为超时
*/ */
private final long expireAt; private volatile long expireAt;
/** /**
* 构造函数 * 构造函数
@ -60,6 +60,10 @@ public class SseEmitterDelayed implements Delayed {
this.expireAt = Instant.now().toEpochMilli() + unit.toMillis(delay); this.expireAt = Instant.now().toEpochMilli() + unit.toMillis(delay);
} }
public void renew(long delay, TimeUnit unit) {
this.expireAt = System.currentTimeMillis() + unit.toMillis(delay);
}
/** /**
* 获取剩余延迟时间 * 获取剩余延迟时间
* *

View File

@ -1,16 +1,22 @@
package org.dromara.common.sse.core; package org.dromara.common.sse.core;
import cn.hutool.core.map.MapUtil; import cn.hutool.core.map.MapUtil;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.dromara.common.core.utils.SpringUtils; import org.dromara.common.core.utils.SpringUtils;
import org.dromara.common.redis.utils.RedisUtils; import org.dromara.common.redis.utils.RedisUtils;
import org.dromara.common.sse.config.SseProperties;
import org.dromara.common.sse.dto.SseMessageDto; import org.dromara.common.sse.dto.SseMessageDto;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException; import java.io.IOException;
import java.lang.ref.WeakReference; import java.lang.ref.WeakReference;
import java.util.Map; import java.util.Map;
import java.util.concurrent.*; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.DelayQueue;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.function.Consumer; import java.util.function.Consumer;
/** /**
@ -33,28 +39,51 @@ public class SseEmitterManager {
*/ */
private static final DelayQueue<SseEmitterDelayed> DELAY_QUEUE = new DelayQueue<>(); private static final DelayQueue<SseEmitterDelayed> DELAY_QUEUE = new DelayQueue<>();
/**
* 心跳事件
*/
private static final SseEmitter.SseEventBuilder PING_EVENT = SseEmitter.event().comment("ping");
private final SseProperties PROPERTIES = SpringUtils.getBean(SseProperties.class);
/** /**
* 清理线程池 * 清理线程池
*/ */
private static final ScheduledExecutorService CLEANER = private final ScheduledExecutorService CLEANER =
Executors.newSingleThreadScheduledExecutor(r -> { Executors.newSingleThreadScheduledExecutor(r -> {
Thread t = new Thread(r, "SSE-Cleaner"); Thread t = new Thread(r, "SSE-Cleaner");
t.setDaemon(true); t.setDaemon(true);
return t; return t;
}); });
static { @PostConstruct
// 每分钟处理到期任务 public void init() {
CLEANER.scheduleWithFixedDelay(SseEmitterManager::processDelayQueue, 1, 1, TimeUnit.MINUTES); log.info("SSE 管理器启动 -> 检查间隔: {} 秒, 心跳间隔: {} 秒",
PROPERTIES.getCheckIntervalSeconds(), PROPERTIES.getHeartbeatSeconds());
CLEANER.scheduleWithFixedDelay(
this::processDelayQueue,
PROPERTIES.getCheckIntervalSeconds(),
PROPERTIES.getCheckIntervalSeconds(),
PROPERTIES.getUnit()
);
}
@PreDestroy
public void destroy() {
CLEANER.shutdownNow();
USER_TOKEN_EMITTERS.clear();
DELAY_QUEUE.clear();
log.info("SSE 管理器已关闭");
} }
/** /**
* 处理延迟队列中的到期 SSE 任务 * 处理延迟队列中的到期 SSE 任务
* 作用 * 作用
* 1. 移除已经失效的 SSEEmitter * 1. 移除已经失效的 SSEEmitter
* 2. 对存活的 SSEEmitter 延迟续期保证心跳检测 * 2. 对存活的 SSEEmitter 延迟续期保证心跳检测
*/ */
private static void processDelayQueue() { private void processDelayQueue() {
try { try {
// 获取当前时间戳用于判断任务是否到期 // 获取当前时间戳用于判断任务是否到期
long now = System.currentTimeMillis(); long now = System.currentTimeMillis();
@ -71,11 +100,11 @@ public class SseEmitterManager {
if (emitter == null || isEmitterDead(emitter)) { if (emitter == null || isEmitterDead(emitter)) {
// Emitter 已被 GC 或已关闭断开连接并从管理器移除 // Emitter 已被 GC 或已关闭断开连接并从管理器移除
SpringUtils.getBean(SseEmitterManager.class).disconnect(task.getUserId(), task.getToken()); this.disconnect(task.getUserId(), task.getToken());
} else { } else {
// Emitter 仍然存活延迟续期 // 直接更新到期时间放回队列
// 5 分钟后再检查该连接避免频繁触发 task.renew(PROPERTIES.getHeartbeatSeconds(), PROPERTIES.getUnit());
DELAY_QUEUE.offer(new SseEmitterDelayed(task.getUserId(), task.getToken(), emitter, 5, TimeUnit.MINUTES)); DELAY_QUEUE.offer(task);
} }
} }
} catch (Exception e) { } catch (Exception e) {
@ -83,7 +112,6 @@ public class SseEmitterManager {
} }
} }
/** /**
* 建立与指定用户的 SSE 连接 * 建立与指定用户的 SSE 连接
* *
@ -100,6 +128,9 @@ public class SseEmitterManager {
Map<String, WeakReference<SseEmitter>> emitters = Map<String, WeakReference<SseEmitter>> emitters =
USER_TOKEN_EMITTERS.computeIfAbsent(userId, k -> new ConcurrentHashMap<>()); USER_TOKEN_EMITTERS.computeIfAbsent(userId, k -> new ConcurrentHashMap<>());
// 如果已有旧连接则先断开
emitters.remove(token);
// 创建一个新的 SseEmitter 实例超时时间设置为一天 避免连接之后直接关闭浏览器导致连接停滞 // 创建一个新的 SseEmitter 实例超时时间设置为一天 避免连接之后直接关闭浏览器导致连接停滞
SseEmitter emitter = new SseEmitter(86400000L); SseEmitter emitter = new SseEmitter(86400000L);
@ -112,7 +143,8 @@ public class SseEmitterManager {
emitter.onError(e -> removeTask.run()); emitter.onError(e -> removeTask.run());
// 延迟清理 // 延迟清理
DELAY_QUEUE.offer(new SseEmitterDelayed(userId, token, emitter, 5, TimeUnit.MINUTES)); DELAY_QUEUE.offer(new SseEmitterDelayed(userId, token, emitter,
PROPERTIES.getHeartbeatSeconds(), PROPERTIES.getUnit()));
try { try {
// 向客户端发送一条连接成功的事件 // 向客户端发送一条连接成功的事件
@ -155,7 +187,6 @@ public class SseEmitterManager {
if (emitters.isEmpty()) { if (emitters.isEmpty()) {
USER_TOKEN_EMITTERS.remove(userId); USER_TOKEN_EMITTERS.remove(userId);
} }
log.debug("SSE连接移除并断开 userId={}, token={}", userId, token); log.debug("SSE连接移除并断开 userId={}, token={}", userId, token);
} }
@ -164,7 +195,7 @@ public class SseEmitterManager {
*/ */
private static boolean isEmitterDead(SseEmitter emitter) { private static boolean isEmitterDead(SseEmitter emitter) {
try { try {
emitter.send(SseEmitter.event().comment("ping")); emitter.send(PING_EVENT);
return false; return false;
} catch (Exception e) { } catch (Exception e) {
return true; return true;