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

Spring WebFlux中SSE端点Flux异常无法传递至客户端的问题

Spring WebFlux SSE端点异常无法传递至客户端的问题

在基于Spring Boot 3.1.5的Spring WebFlux项目中,我通过Flux.create和sink实现SSE端点。当初始化完成后抛出异常并发送时,连接直接断开,异常并未传递至客户端,但我期望客户端Flux能接收该异常。

调试发现问题在于Spring判定连接已提交,但我认为仍可通过开放连接传递异常。日志显示:
Error [java.lang.RuntimeException: Wifi cable not found 2023-11-08T07:45:36.574336976Z] for HTTP GET "/api/fail/hello-flux-later", but ServerHttpResponse already committed (200 OK)

我不想通过捕获异常并在响应对象中添加错误字段,或包装所有响应实体的方式解决,希望找到无需显式处理异常的通用方案。或许可自定义SSE的error字段,但该字段未被标准规范定义,这可能是Spring直接中断连接的原因?

(注:我知晓Flux.create不符合响应式理念,此处用于桥接场景,非问题核心)

服务端代码

public record SimpleResponse(String content) {}

@GetMapping(path = "hello-flux-later", produces = TEXT_EVENT_STREAM_VALUE)
Flux<SimpleResponse> helloFluxLater() {
    return Flux.create(sink -> {

        final Thread creator = new Thread(() -> {
            LOGGER.info("Starting hello flux thread, failing after 3");

            for (int i = 0; i < 5; ++i) {
                if (i == 3) {

                    // FIXME: not received by the subscriber
                    LOGGER.info("Sending {} error", i);
                    sink.error(new RuntimeException("Wifi cable not found " + Instant.now()));
                } else {

                    LOGGER.info("Sending {}", i);
                    sink.next(new SimpleResponse("hello " + Instant.now()));
                }

                try {
                    Thread.sleep(1000);
                } catch (final InterruptedException e) {
                    LOGGER.info("Interrupted at {}", i);
                    return;
                }
            }

            LOGGER.info("Sending done, finishing sink and thread");
            sink.complete();
        });
        creator.setDaemon(true);
        creator.start();

        sink.onDispose(() -> {
            LOGGER.info("Flux disposed, stopping thread");
            creator.interrupt();
        });
    });
}

客户端请求日志

http -v --stream localhost:8080/api/fail/hello-flux-later                                                      
GET /api/fail/hello-flux-later HTTP/1.1
Accept: */*
Accept-Encoding: gzip, deflate
Connection: keep-alive
Host: localhost:8080
User-Agent: HTTPie/2.2.0



HTTP/1.1 200 OK
Content-Type: text/event-stream;charset=UTF-8
transfer-encoding: chunked

data:{"content":"hello 2023-11-08T07:45:33.568766666Z"}

data:{"content":"hello 2023-11-08T07:45:34.569467283Z"}

data:{"content":"hello 2023-11-08T07:45:35.573303790Z"}


http: error: ChunkedEncodingError: ("Connection broken: InvalidChunkLength(got length b'', 0 bytes read)", InvalidChunkLength(got length b'', 0 bytes read))

服务端日志

2023-11-08T08:45:33.563+01:00 DEBUG 10893 --- [or-http-epoll-9] r.n.http.server.HttpServerOperations     : [11a660e4, L:/[0:0:0:0:0:0:0:1%0]:8080 - R:/[0:0:0:0:0:0:0:1%0]:44362] New http connection, requesting read
2023-11-08T08:45:33.563+01:00 DEBUG 10893 --- [or-http-epoll-9] r.netty.transport.TransportConfig        : [11a660e4, L:/[0:0:0:0:0:0:0:1%0]:8080 - R:/[0:0:0:0:0:0:0:1%0]:44362] Initialized pipeline DefaultChannelPipeline{(reactor.left.httpCodec = io.netty.handler.codec.http.HttpServerCodec), (reactor.left.httpTrafficHandler = reactor.netty.http.server.HttpTrafficHandler), (reactor.right.reactiveBridge = reactor.netty.channel.ChannelOperationsHandler)}
2023-11-08T08:45:33.566+01:00 DEBUG 10893 --- [or-http-epoll-9] r.n.http.server.HttpServerOperations     : [11a660e4, L:/[0:0:0:0:0:0:0:1%0]:8080 - R:/[0:0:0:0:0:0:0:1%0]:44362] Increasing pending responses, now 1
2023-11-08T08:45:33.567+01:00 DEBUG 10893 --- [or-http-epoll-9] reactor.netty.http.server.HttpServer     : [11a660e4-1, L:/[0:0:0:0:0:0:0:1%0]:8080 - R:/[0:0:0:0:0:0:0:1%0]:44362] Handler is being applied: org.springframework.http.server.reactive.ReactorHttpHandlerAdapter@7ab3fac
2023-11-08T08:45:33.567+01:00 DEBUG 10893 --- [or-http-epoll-9] o.s.w.s.adapter.HttpWebHandlerAdapter    : [11a660e4-8] HTTP GET "/api/fail/hello-flux-later"
2023-11-08T08:45:33.567+01:00 DEBUG 10893 --- [or-http-epoll-9] s.w.r.r.m.a.RequestMappingHandlerMapping : [11a660e4-8] Mapped to de.gsi.bpeter.experiment.springwebfluxplain.FailingController#helloFluxLater()
2023-11-08T08:45:33.568+01:00 DEBUG 10893 --- [or-http-epoll-9] o.s.w.r.r.m.a.ResponseBodyResultHandler  : [11a660e4-8] Using 'text/event-stream' given [*/*] and supported [text/event-stream]
2023-11-08T08:45:33.568+01:00 DEBUG 10893 --- [or-http-epoll-9] o.s.w.r.r.m.a.ResponseBodyResultHandler  : [11a660e4-8] 0..N [de.gsi.bpeter.experiment.springwebfluxplain.SimpleResponse]
2023-11-08T08:45:33.568+01:00 DEBUG 10893 --- [or-http-epoll-9] reactor.netty.ReactorNetty               : [11a660e4-1, L:/[0:0:0:0:0:0:0:1%0]:8080 - R:/[0:0:0:0:0:0:0:1%0]:44362] Non Removed handler: reactor.left.readTimeoutHandler, context: null, pipeline: DefaultChannelPipeline{(reactor.left.httpCodec = io.netty.handler.codec.http.HttpServerCodec), (reactor.left.httpTrafficHandler = reactor.netty.http.server.HttpTrafficHandler), (reactor.right.reactiveBridge = reactor.netty.channel.ChannelOperationsHandler)}
2023-11-08T08:45:33.568+01:00  INFO 10893 --- [      Thread-15] d.g.b.e.s.SpringWebfluxPlainService      : Starting hello flux thread, failing after 3
2023-11-08T08:45:33.568+01:00  INFO 10893 --- [      Thread-15] d.g.b.e.s.SpringWebfluxPlainService      : Sending 0
2023-11-08T08:45:33.568+01:00 DEBUG 10893 --- [      Thread-15] d.g.b.e.s.SimpleResponse                 : Creating response SimpleResponse[content=hello 2023-11-08T07:45:33.568766666Z]
2023-11-08T08:45:33.569+01:00 DEBUG 10893 --- [or-http-epoll-9] o.s.http.codec.json.Jackson2JsonEncoder  : [11a660e4-8] Encoding [SimpleResponse[content=hello 2023-11-08T07:45:33.568766666Z]]
2023-11-08T08:45:34.569+01:00  INFO 10893 --- [      Thread-15] d.g.b.e.s.SpringWebfluxPlainService      : Sending 1
2023-11-08T08:45:34.569+01:00 DEBUG 10893 --- [      Thread-15] d.g.b.e.s.SimpleResponse                 : Creating response SimpleResponse[content=hello 2023-11-08T07:45:34.569467283Z]
2023-11-08T08:45:34.569+01:00 DEBUG 10893 --- [      Thread-15] o.s.http.codec.json.Jackson2JsonEncoder  : [11a660e4-8] Encoding [SimpleResponse[content=hello 2023-11-08T07:45:34.569467283Z]]
2023-11-08T08:45:35.573+01:00  INFO 10893 --- [      Thread-15] d.g.b.e.s.SpringWebfluxPlainService      : Sending 2
2023-11-08T08:45:35.573+01:00 DEBUG 10893 --- [      Thread-15] d.g.b.e.s.SimpleResponse                 : Creating response SimpleResponse[content=hello 2023-11-08T07:45:35.573303790Z]
2023-11-08T08:45:35.573+01:00 DEBUG 10893 --- [      Thread-15] o.s.http.codec.json.Jackson2JsonEncoder  : [11a660e4-8] Encoding [SimpleResponse[content=hello 2023-11-08T07:45:35.573303790Z]]
2023-11-08T08:45:36.574+01:00  INFO 10893 --- [      Thread-15] d.g.b.e.s.SpringWebfluxPlainService      : Sending 3 error
2023-11-08T08:45:36.574+01:00 DEBUG 10893 --- [      Thread-15] d.g.b.e.s.SimpleResponse                 : Creating exception java.lang.RuntimeException: Wifi cable not found 2023-11-08T07:45:36.574336976Z
2023-11-08T08:45:36.574+01:00  INFO 10893 --- [      Thread-15] d.g.b.e.s.SpringWebfluxPlainService      : Flux disposed, stopping thread
2023-11-08T08:45:36.574+01:00  INFO 10893 --- [      Thread-15] d.g.b.e.s.SpringWebfluxPlainService      : Interrupted at 3
2023-11-08T08:45:36.574+01:00 ERROR 10893 --- [or-http-epoll-9] o.s.w.s.adapter.HttpWebHandlerAdapter    : [11a660e4-8] Error [java.lang.RuntimeException: Wifi cable not found 2023-11-08T07:45:36.574336976Z] for HTTP GET "/api/fail/hello-flux-later", but ServerHttpResponse already committed (200 OK)
2023-11-08T08:45:36.575+01:00 ERROR 10893 --- [or-http-epoll-9] r.n.http.server.HttpServerOperations     : [11a660e4-1, L:/[0:0:0:0:0:0:0:1%0]:8080 - R:/[0:0:0:0:0:0:0:1%0]:44362] Error finishing response. Closing connection

java.lang.RuntimeException: Wifi cable not found 2023-11-08T07:45:36.574336976Z
    at de.gsi.bpeter.experiment.springwebfluxplain.SimpleResponse.createException(SimpleResponse.java:19) ~[classes/:na]
    Suppressed: reactor.core.publisher.FluxOnAssembly$OnAssemblyException: 
Error has been observed at the following site(s):
    *__checkpoint ⇢ Handler de.gsi.bpeter.experiment.springwebfluxplain.FailingController#helloFluxLater() [DispatcherHandler]
    *__checkpoint ⇢ HTTP GET "/api/fail/hello-flux-later" [ExceptionHandlingWebHandler]
Original Stack Trace:
        at de.gsi.bpeter.experiment.springwebfluxplain.SimpleResponse.createException(SimpleResponse.java:19) ~[classes/:na]
        at de.gsi.bpeter.experiment.springwebfluxplain.FailingController.lambda$2(FailingController.java:60) ~[classes/:na]
        at java.base/java.lang.Thread.run(Thread.java:833) ~[na:na]

2023-11-08T08:45:36.575+01:00 DEBUG 10893 --- [or-http-epoll-9] r.netty.channel.ChannelOperations        : [11a660e4-1, L:/[0:0:0:0:0:0:0:1%0]:8080 ! R:/[0:0:0:0:0:0:0:1%0]:44362] [HttpServer] Channel inbound receiver cancelled (channel disconnected).

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 14:55:55