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

为何fetch event source创建后即刻关闭?Angular+Spring Boot SSE问题

SSE连接创建后立即关闭的解决方案

我尝试使用微软fetch-event-source包,从Spring Boot服务器向Angular客户端推送即时及定时通知,但Event Source创建后不久便自动关闭,求解决方案。


Angular 服务类

constructor(private http:HttpClient, 
    private authService:AuthService, private zone:NgZone) {

      const currentUser = toSignal(this.authService.loggedInUserObs$);

      effect((onCleanup) => {

        const user = currentUser?.();

        if(user && this.isAdminUser(user)){  

          this.connectToNotifications(user)
        } else {

          this.disconnectFromSSE();
        }

        onCleanup(() => {
          this.disconnectFromSSE();
        }
        )
      }, {allowSignalWrites:true});
   
      
      
     }

    //  disconnect from SSE notifications
  disconnectFromSSE() {
    
   this.abortController?.abort();
   this.abortController = null;
   this.connectionState.set('disconnected');
  }
  private async connectToNotifications(user: User) {

 
    //this.disconnectFromSSE();
    this.connectionState.set('connecting');

    this.abortController = new AbortController();

    const token = user.accessToken;

   

    try{

      await fetchEventSource(`${this.connectionUrl}?_xxid=${user.id}`, {
        method: 'GET',
        // credentials: 'include',
        headers: {
          'Authorization': `Bearer ${token}`,
          'Accept': 'text/event-stream',
          'Cache-Control': 'no-cache',
          'Connection': 'keep-alive',
          

        },

       
        signal: this.abortController.signal,
        onopen: async (response) => {

          this.zone.run(() => {

            if(response.ok && response.status === 200){

              console.log('200 ok response')

              this.connectionState.set('connected');

            this.retryCount = 0;

            return;

            
            }else if(response.status >= 400){
              console.log('4xx error response')
              this.connectionState.set('error');
            }
            

            throw new Error(`Connection failed: ${response.status} ${response.statusText}`);
          });

        },

        onmessage:(event:EventSourceMessage) => {
          
          this.zone.run(() => {

            try{

              if(event.event==='responseUpdate'){


                const data:AssessmentResponseRecord = JSON.parse(event.data);

                console.log(JSON.stringify(data, null,1))

               const index = this.notifications().findIndex((n) => n.topic === data.topic && n.instructorId === data.instructorId && n.postedOn === data.postedOn);
               if(index === -1){
   
                this.notifications.update(prev => [...prev, data]);
               

               }
              }

            } catch (error) {
              console.error('Error processing event message:', error);
            }
          });
        },

        onerror:(err) => {

          this.zone.run(() => {

            console.error('SSE error', err);
            this.connectionState.set('error');

            if(err.name === 'AbortError') return;

              this.attemptReconnection(user);
            
          });
        },

        onclose: () => {
          this.zone.run(() => {
            console.log('SSE connection closed');
            this.connectionState.set('disconnected');

          
            
          });
        }
      })
    }catch(err){

      this.zone.run(() => {
        this.connectionState.set('error');
        console.error('SSE connection failed:', err);
        
      //  this.attemptReconnection(user);
      })
    }

  }

  private attemptReconnection(user: User) {
    if (this.retryCount < this.maxRetries) {

      const delay = Math.min(this.baseDelay * 2 ** this.retryCount, 30000);

      this.retryCount++;

      setTimeout(() => this.connectToNotifications(user), delay);
    }
  }

Spring Boot 连接服务

@Service
public class ConnectorService{
private Map<Integer, SseEmitter> connectors = new ConcurrentHashMap<>();
    
    
    @Override
    public SseEmitter establishConnection(Integer instructorId) {
      
//      get previous connector if there was
        SseEmitter emitter = connectors.get(instructorId);

//      create new connector if there was non previously
        if(emitter == null) {
            
            emitter = new SseEmitter(0L);
            
            connectors.put(instructorId, emitter);
            
            System.out.println("new connection made");
        }else {
            
            System.out.println("previously connected");
        }
        heartbeat(emitter);
        
        emitter.onCompletion(() => connectors.remove(instructorId));
        
        emitter.onTimeout(() => connectors.remove(instructorId));
        
    
    
        
        return emitter;
    }
}

定时通知服务

@Service
public class ScheduledNotificationService{
@Scheduled(fixedRate = 100000) //schedule notifications at 20sec interval
    public void scheduleFixedTimeNotification() {
        
    
        
        RLock rlock = redissonClient.getLock("notification:lock");
        try {
            
            if(rlock.tryLock(5, 15, TimeUnit.SECONDS)) {
                
                
                
                try {
                    Set<String> keys = cachingKeyUtils.getAllCacheKeys(RedisValues.ASSESSMENT_RESPONSE_NOTIFICATION);
                    
                    Cache cache = cacheManager.getCache(RedisValues.ASSESSMENT_RESPONSE_NOTIFICATION);
                    
                    
                    
                    for(String key: keys) {
                        
                        
                        try {
                            
                            
                            Cache.ValueWrapper wrapper = cache.get(key);
                            
                            if(wrapper != null) {
                                
//                              
                                List<AssessmentResponseRecord> notifications = mapper.convertValue(wrapper.get(),
                                        
                                        new TypeReference<List<AssessmentResponseRecord>>() {}
                                        );
                                
                                broadcaster.broadcastPreviousNotifications(notifications);
                                
                                
                                
                            }
                            
    
                        } catch (Exception e) {
                            System.err.println(e);
                        }
                    }
                } finally  {
                    rlock.unlock();
                }
            }
        } catch (Exception e) {
            Thread.currentThread().interrupt();
            System.err.println(e);
        }
        
        
        }
        

问题排查与解决方案

1. 补充服务器端心跳机制

代码中调用了heartbeat(emitter)但未实现,SSE连接长时间无数据时,浏览器或代理会主动断开。添加心跳实现:

private void heartbeat(SseEmitter emitter) {
    ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
    scheduler.scheduleAtFixedRate(() -> {
        try {
            // 发送空心跳事件维持连接
            emitter.send(SseEmitter.event().name("heartbeat").data(""));
        } catch (IOException e) {
            scheduler.shutdown();
            connectors.remove(emitter);
        }
    }, 0, 15, TimeUnit.SECONDS); // 每15秒发送一次心跳
}

2. 完善Spring Boot控制器接口

缺少暴露给前端的Controller,添加接口确保正确返回SseEmitter:

@RestController
@RequestMapping("/api/notifications")
public class NotificationController {

    @Autowired
    private ConnectorService connectorService;

    @GetMapping("/stream")
    public SseEmitter streamNotifications(@RequestParam("_xxid") Integer instructorId) {
        return connectorService.establishConnection(instructorId);
    }
}

同时确保响应头正确:Content-Type: text/event-stream、Cache-Control: no-cache,避免拦截器修改这些头。

3. 修复Angular端连接配置

  • 移除手动设置的Connection: keep-alive,fetch-event-source会自动处理
  • 取消注释connectToNotifications开头的this.disconnectFromSSE(),避免重复连接导致旧连接被终止
  • 在onopen中打印response.headers.get('Content-Type'),确认返回类型为text/event-stream

4. 调整服务器超时设置

在application.properties中添加配置,避免容器强制断开长连接:

server.tomcat.connection-timeout=0
server.servlet.session.timeout=3600s

5. 排查日志与控制台

  • 查看浏览器控制台是否有CORS、4xx/5xx错误
  • 检查服务器日志,确认连接建立、心跳发送、消息推送是否正常

内容的提问来源于stack exchange,提问作者Anyanwu Chinedu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 08:50:00