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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 10:35:21