Spring SSE(SseEmitter):Flutter客户端仅超时后接收事件问题
Spring Boot SSE事件无法实时推送到Flutter客户端
环境与实现背景
- 后端:Spring Boot 2.7.0 + Spring 5.3.20,基于Spring MVC的Server-Sent-Events(SSE)实现服务端到客户端异步通信
- 客户端:Flutter,通过http包订阅SSE流
服务端代码实现
StreamingController
@RestController @RequestMapping(value = ApiConstants.STREAM_API) @Validated public class StreamingController { private final StreamingService service; @GetMapping(value = "/subscribe", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter subscribe() { SseEmitter emitter = new SseEmitter(120000L); //Timeout 2 Minutes try { emitter.send(SseEmitter.event().name("INIT")); } catch (Exception ex) { ex.printStackTrace(); } emitter.onCompletion(() -> service.getEmitters().remove(emitter)); emitter.onTimeout(()-> service.getEmitters().remove(emitter)); emitter.onError((ex)-> service.getEmitters().remove(emitter)); service.getEmitters().add(emitter); return emitter; } }
StreamingService
@Service @Transactional public class StreamingService { @Getter @Setter private List<SseEmitter> emitters = new CopyOnWriteArrayList<>(); }
EntityListener(实体持久化后发送事件)
@Service @Transactional public class EntityListener { private final StreamingService streamingService; @PostPersist protected void afterCreate(final Entity createdEntity) { List<SseEmitter> emitters = new ArrayList<>(streamingService.getEmitters()); for (SseEmitter emitter : emitters) { try { SseEmitter.SseEventBuilder event = SseEmitter.event() .data("Last Score" + createdEntity.getScore()) .id(String.valueOf(createdEntity.getId())) .name("Event Name"); emitter.send(event); } catch (Exception ex) { emitter.completeWithError(ex); streamingService.getEmitters().remove(emitter); } } } }
客户端代码实现(Flutter)
print("Subscribing.."); Future<http.StreamedResponse>? response; try { final _client = http.Client(); var request = http.Request("GET", Uri.parse('http://localhost:5555/stream/subscribe')); Map<String, String> headers = {}; headers.addAll(service.header); headers["Authorization"] = __token!; headers["Cache-Control"] = "no-cache"; headers["Accept"] = "text/event-stream"; request.headers.addAll(headers); response = _client.send(request); print("Subscribed!"); } catch (e) { print("Caught $e"); } response?.asStream().listen((streamedResponse) { print("Received streamedResponse.statusCode:${streamedResponse.statusCode}"); streamedResponse.stream.listen((data) { print("Received data:${utf8.decode(data)}"); }); });
当前问题
调用emitter.send()方法后事件不会立即推送,所有事件会在SseEmitter超时后同时到达客户端,无法实现Flutter监听器实时接收事件的需求。
解决方案
1. 脱离事务上下文发送事件
EntityListener的@PostPersist方法运行在事务上下文内,SSE事件的发送操作会被事务提交机制延迟,直到事务完成才会输出响应。需将事件发送逻辑异步执行,脱离当前事务:
修改EntityListener代码:
@Service public class EntityListener { private final StreamingService streamingService; private final TaskExecutor taskExecutor; public EntityListener(StreamingService streamingService, TaskExecutor taskExecutor) { this.streamingService = streamingService; this.taskExecutor = taskExecutor; } @PostPersist protected void afterCreate(final Entity createdEntity) { // 异步执行事件发送,避开事务上下文 taskExecutor.execute(() -> { List<SseEmitter> emitters = new ArrayList<>(streamingService.getEmitters()); for (SseEmitter emitter : emitters) { try { SseEmitter.SseEventBuilder event = SseEmitter.event() .data("Last Score" + createdEntity.getScore()) .id(String.valueOf(createdEntity.getId())) .name("Event Name"); emitter.send(event); } catch (Exception ex) { emitter.completeWithError(ex); streamingService.getEmitters().remove(emitter); } } }); } }
添加Spring异步配置类:
@Configuration @EnableAsync public class AsyncConfig { @Bean public TaskExecutor taskExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(5); executor.setMaxPoolSize(10); executor.setQueueCapacity(25); return executor; } }
2. 禁用响应缓冲
Spring MVC默认可能对响应进行缓冲,需手动设置响应头确保事件实时输出:
修改StreamingController的subscribe方法:
@GetMapping(value = "/subscribe", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter subscribe(HttpServletResponse response) { // 禁用响应缓冲,确保SSE事件实时推送 response.setHeader("Cache-Control", "no-cache"); response.setHeader("X-Accel-Buffering", "no"); // 适配Nginx等反向代理场景 response.setContentType(MediaType.TEXT_EVENT_STREAM_VALUE); SseEmitter emitter = new SseEmitter(120000L); //Timeout 2 Minutes try { emitter.send(SseEmitter.event().name("INIT")); } catch (Exception ex) { ex.printStackTrace(); } emitter.onCompletion(() -> service.getEmitters().remove(emitter)); emitter.onTimeout(()-> service.getEmitters().remove(emitter)); emitter.onError((ex)-> service.getEmitters().remove(emitter)); service.getEmitters().add(emitter); return emitter; }
3. 优化Flutter客户端流处理
简化客户端的流监听逻辑,直接处理send请求返回的Future,避免不必要的流转换:
print("Subscribing.."); try { final _client = http.Client(); var request = http.Request("GET", Uri.parse('http://localhost:5555/stream/subscribe')); Map<String, String> headers = {}; headers.addAll(service.header); headers["Authorization"] = __token!; headers["Cache-Control"] = "no-cache"; headers["Accept"] = "text/event-stream"; request.headers.addAll(headers); final streamedResponse = await _client.send(request); print("Subscribed! Status code: ${streamedResponse.statusCode}"); streamedResponse.stream.listen((data) { print("Received data: ${utf8.decode(data)}"); }, onError: (e) { print("Stream error: $e"); }, onDone: () { print("Stream closed"); }); } catch (e) { print("Caught $e"); }
内容的提问来源于stack exchange,提问作者HLA
相关产品推荐
相关产品推荐

