You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

React页面切换时Server-Sent Event连接未关闭问题求助

问题:SSE连接重复创建导致推送间隔异常

问题描述

基于SpringBoot+React的全栈应用,后端通过SSE接口推送包含time、pushcode、successes、failures字段的Java对象栈数据,前端用EventSource建立连接。正常页面切换时功能正常,但页面刷新或跳转返回后,出现多个SSE连接同时推送的情况,原本4秒的间隔被打乱,数据重复推送。

相关代码

前端组件代码

const WorkingSmsesEvent = () => {
  const executedSmsListener = useStore((state) => state.executedSmsListener)
  const setExecutedSmsListener = useStore((state) => state.setExecutedSmsListener)
  
  const [data, setData] = useState(null)
  let eventSource = undefined;

  const handleMessage = (event) => {
    const msg = JSON.parse(event.data);
    console.log("addEventListener Progress:", event.target.readyState);
    setData(msg)
  }

  useEffect(() => {
    if (!executedSmsListener) {
      eventSource = new EventSource("http://localhost:9090/api/dashboard/subscribe", { withCredentials: true});
      
      window.onbeforeunload = function() {
        eventSource.close();
        eventSource.removeEventListener("Progress", handleMessage)
        setExecutedSmsListener(false);
        return true;
      };

      eventSource.addEventListener("Progress", handleMessage);

      eventSource.onerror = (event) => {
        if (event.target.readyState === EventSource.CLOSED) {
          console.log('SSE closed (' + event.target.readyState + ')')
        }
        eventSource.close();
      }

      eventSource.onopen = (event) => {
        setExecutedSmsListener(true);
      }
      
      return () => {
        setExecutedSmsListener(false);
        eventSource.close();   
        eventSource.removeEventListener("Progress", handleMessage)
      }
    }
  }, [])

  return (
    <div className='flex flex-col'>
        <div className='flex flex-wrap'>
        {
            data?.map((item, id) => (
            <div key={id}>
                <Cards extend={false} data={item}/>
            </div>
            ))
        }
        </div>
    </div>
  );
}

Zustand状态管理代码

import { create } from 'zustand'

export const useStore = create((set) => ({
  executedSmsListener: false,
  setExecutedSmsListener: executedSmsListener => set((state) => ({ executedSmsListener }))
}))

后端Controller代码

@GetMapping("/subscribe")
public ResponseEntity<SseEmitter> subscribeSSE() throws InterruptedException, IOException {
    final SseEmitter emitter = new SseEmitter((long)-1);
    service.addEmitter(emitter);
    service.subscribeExecutedSms();
    emitter.onCompletion(() -> service.removeEmitter(emitter));
    emitter.onTimeout(() -> service.removeEmitter(emitter));
    return new ResponseEntity<>(emitter, HttpStatus.OK);
}

后端Service代码

public void subscribeExecutedSms() throws IOException {
    List<SseEmitter> deadEmitters = new ArrayList<>();
    Random num = new Random();
    
    scheduled = taskScheduler.scheduleAtFixedRate(() => {
        ExecutedSms executed = new ExecutedSms(TIME_FORMATTER.format(new Date()), num.nextInt(20) + 2, num.nextInt(3), num.nextInt(10) + 115);
        execSmsList.add(executed);
        execSmsEmitters.forEach(emitter => {
            try {
                emitter.send(SseEmitter.event().name("Progress")
                        .data(execSmsList));
            } catch (Exception e) {
                deadEmitters.add(emitter);
            }
        });
        if (execSmsList.size() == 8) {
            execSmsList.clear();
        }
        execSmsEmitters.removeAll(deadEmitters);
    }, (long)4000);
}

问题排查思路

前端问题点

  1. window.onbeforeunload不可靠:页面刷新/跳转时,浏览器可能在回调执行完成前就卸载页面,导致连接未关闭、状态未重置。
  2. 内存状态丢失:Zustand的executedSmsListener存在内存中,页面刷新后会重置为false,触发新连接创建。
  3. useEffect依赖缺失:依赖数组为空,executedSmsListener状态变化时无法重新触发连接检查逻辑。
  4. 连接实例未持久化:eventSource定义在组件内部,刷新后旧连接可能未被后端清理,新连接又建立。

后端问题点

  1. 重复创建定时任务:每次前端请求SSE接口,都会调用subscribeExecutedSms新建定时任务,导致多个任务同时推送数据。
  2. 定时任务与Emitter耦合:定时任务未与Emitter生命周期绑定,前端关闭连接后,后端任务仍在运行并推送数据。
  3. 非线程安全集合:execSmsList和execSmsEmitters未使用线程安全容器,并发场景下可能出现数据不一致。

解决方案

前端优化

  1. 移除window.onbeforeunload:完全依赖React的useEffect清理函数,组件卸载时自动执行,可靠性更高。
  2. 用useRef保存连接实例:确保能跟踪当前活跃连接,避免组件重渲染时丢失引用。
  3. 补全useEffect依赖:添加executedSmsListener和setExecutedSmsListener,确保状态变化时重新检查连接状态。

修改后的前端代码:

const WorkingSmsesEvent = () => {
  const executedSmsListener = useStore((state) => state.executedSmsListener)
  const setExecutedSmsListener = useStore((state) => state.setExecutedSmsListener)
  
  const [data, setData] = useState(null)
  const eventSourceRef = useRef(null); // 用useRef持久化连接实例

  const handleMessage = (event) => {
    const msg = JSON.parse(event.data);
    console.log("addEventListener Progress:", event.target.readyState);
    setData(msg)
  }

  useEffect(() => {
    // 已有活跃连接或已标记监听,直接返回
    if (executedSmsListener || eventSourceRef.current) {
      return;
    }

    const eventSource = new EventSource("http://localhost:9090/api/dashboard/subscribe", { withCredentials: true});
    eventSourceRef.current = eventSource;

    eventSource.addEventListener("Progress", handleMessage);

    eventSource.onerror = (event) => {
      console.log('SSE error:', event.target.readyState);
      eventSource.close();
      eventSourceRef.current = null;
      setExecutedSmsListener(false);
    }

    eventSource.onopen = (event) => {
      setExecutedSmsListener(true);
    }

    // 组件卸载时的清理逻辑
    return () => {
      if (eventSourceRef.current) {
        eventSourceRef.current.close();
        eventSourceRef.current.removeEventListener("Progress", handleMessage);
        eventSourceRef.current = null;
      }
      setExecutedSmsListener(false);
    }
  }, [executedSmsListener, setExecutedSmsListener]) // 补全依赖

  return (
    <div className='flex flex-col'>
        <div className='flex flex-wrap'>
        {
            data?.map((item, id) => (
            <div key={id}>
                <Cards extend={false} data={item}/>
            </div>
            ))
        }
        </div>
    </div>
  );
}

后端优化

  1. 全局单一定时任务:仅在有活跃Emitter时启动定时任务,无连接时自动停止,避免重复创建。
  2. 线程安全容器:用Collections.synchronizedList和Collections.synchronizedSet存储数据和Emitter,避免并发问题。
  3. Emitter生命周期绑定:新连接建立时推送历史数据,连接关闭时自动移除并检查是否需要停止定时任务。

修改后的后端Service代码:

import org.springframework.scheduling.TaskScheduler;
import org.springframework.scheduling.concurrent.ScheduledFuture;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;

import java.io.IOException;
import java.text.SimpleDateFormat;
import java.util.*;
import java.util.concurrent.TimeUnit;

public class SseService {
    private ScheduledFuture<?> scheduledTask;
    private final List<ExecutedSms> execSmsList = Collections.synchronizedList(new ArrayList<>());
    private final Set<SseEmitter> execSmsEmitters = Collections.synchronizedSet(new HashSet<>());
    private final Random num = new Random();
    private final TaskScheduler taskScheduler;
    private static final SimpleDateFormat TIME_FORMATTER = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");

    public SseService(TaskScheduler taskScheduler) {
        this.taskScheduler = taskScheduler;
    }

    public void addEmitter(SseEmitter emitter) {
        execSmsEmitters.add(emitter);
        // 新连接建立时推送当前历史数据
        try {
            emitter.send(SseEmitter.event().name("Progress").data(execSmsList));
        } catch (IOException e) {
            execSmsEmitters.remove(emitter);
        }
        // 无活跃定时任务时启动
        if (scheduledTask == null || scheduledTask.isCancelled()) {
            startScheduledTask();
        }
    }

    public void removeEmitter(SseEmitter emitter) {
        execSmsEmitters.remove(emitter);
        // 无活跃连接时停止定时任务
        if (execSmsEmitters.isEmpty() && scheduledTask != null && !scheduledTask.isCancelled()) {
            scheduledTask.cancel(true);
            scheduledTask = null;
        }
    }

    private void startScheduledTask() {
        scheduledTask = taskScheduler.scheduleAtFixedRate(() -> {
            // 生成新数据
            ExecutedSms executed = new ExecutedSms(TIME_FORMATTER.format(new Date()), num.nextInt(20) + 2, num.nextInt(3), num.nextInt(10) + 115);
            execSmsList.add(executed);
            if (execSmsList.size() == 8) {
                execSmsList.clear();
            }
            // 推送数据给所有活跃Emitter
            List<SseEmitter> deadEmitters = new ArrayList<>();
            execSmsEmitters.forEach(emitter -> {
                try {
                    emitter.send(SseEmitter.event().name("Progress").data(execSmsList));
                } catch (Exception e) {
                    deadEmitters.add(emitter);
                }
            });
            execSmsEmitters.removeAll(deadEmitters);
        }, 4, TimeUnit.SECONDS);
    }

    // 原subscribeExecutedSms改为空实现,逻辑已迁移到startScheduledTask和addEmitter
    public void subscribeExecutedSms() throws IOException {
    }
}

额外注意点

  • 前端:避免在多个组件中重复创建SSE连接,建议将SSE逻辑封装为全局Hook或单例服务。
  • 后端:添加Emitter超时时间(不建议设为-1),避免无效连接长期占用资源;定时任务中捕获所有异常,防止单个Emitter报错导致任务终止。

内容的提问来源于stack exchange,提问作者Shiladitya Thakur

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 10:43:07