diff --git a/src/main/java/cn/lihongjie/coal/spring/config/GlobalExceptionHandler.java b/src/main/java/cn/lihongjie/coal/spring/config/GlobalExceptionHandler.java index fb4b5b2e..cd8dc93a 100644 --- a/src/main/java/cn/lihongjie/coal/spring/config/GlobalExceptionHandler.java +++ b/src/main/java/cn/lihongjie/coal/spring/config/GlobalExceptionHandler.java @@ -6,6 +6,7 @@ import cn.lihongjie.coal.exception.BizException; import jakarta.servlet.http.HttpServletRequest; import jakarta.validation.ConstraintViolation; +import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import org.hibernate.SessionFactory; @@ -21,6 +22,7 @@ import org.springframework.web.bind.annotation.ExceptionHandler; import org.springframework.web.bind.annotation.RestControllerAdvice; import org.springframework.web.method.HandlerMethod; +import java.io.IOException; import java.util.Arrays; import java.util.List; import java.util.Map; @@ -65,6 +67,28 @@ public class GlobalExceptionHandler { return R.fail(ex.getCode(), ex.getMessage()); } + @SneakyThrows + @ExceptionHandler(IOException.class) + public void handleIOException(IOException ex, HttpServletRequest request) { + boolean isSse = "text/event-stream".equals(request.getHeader("Accept")); + String message = ex.getMessage(); + // 忽略常见的客户端断开异常 + if ( isSse && message != null && ( + message.contains("Broken pipe") || + message.contains("Connection reset by peer") || + message.contains("An existing connection was forcibly closed") || + message.contains("你的主机中的软件中止了一个已建立的连接") + )) { + // 可以选择只记录DEBUG日志,或者直接忽略 + // log.debug("SSE client closed connection: {}", message); + log.info("SSE client closed connection: {}", message); + return; + } + // 如果不是上述异常,可以选择抛出或做其他处理 + // log.error("IOException not ignored", ex); + throw ex; // 或返回错误响应 + } + private void logEx(Exception ex, HttpServletRequest request, HandlerMethod handlerMethod) { if (request.getAttribute("__logged") == null) { diff --git a/src/main/java/cn/lihongjie/coal/sse/service/SseService.java b/src/main/java/cn/lihongjie/coal/sse/service/SseService.java index 0292d012..26c75da1 100644 --- a/src/main/java/cn/lihongjie/coal/sse/service/SseService.java +++ b/src/main/java/cn/lihongjie/coal/sse/service/SseService.java @@ -29,6 +29,8 @@ public class SseService { public void register(String topicId, SseEmitter emitter) { + log.info("{} register sse emitter:{}", topicId, emitter); + ScheduledFuture scheduledFuture = taskScheduler.scheduleAtFixedRate( () -> { try { @@ -42,16 +44,11 @@ public class SseService { } catch (IOException e) { log.error("{} send heartbeat error:{}", topicId, e.getMessage()); - try{ - emitter.complete(); - }catch (Exception exception){ - log.error("{} emitter complete error:{}", topicId, exception.getMessage()); - } } }, - Duration.ofMinutes(1)); + Duration.ofSeconds(30)); RTopic topic = redissonClient.getTopic(topicId); int addListener = @@ -83,12 +80,15 @@ public class SseService { emitter.onError( (e) -> { topic.removeListener(addListener); - log.info("{} remove listener:{} onError ", topicId, addListener, e); + log.info("{} remove listener:{} onError {}", topicId, addListener, e.getMessage()); scheduledFuture.cancel(true); + }); + + } @SneakyThrows