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
相关产品推荐
相关产品推荐

