部署到Kubernetes时Spring Webflux的doOnCancel未触发问题
我有一个提供Flux流的Controller端点,本地环境下关闭浏览器标签页时,doOnCancel和doOnTerminate方法能正常触发;但将应用部署到Kubernetes后,这两个方法无法被调用。
Controller代码
@Slf4j @RestController public class TestController { ... @GetMapping(value = "/test", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> testStream() { log.info("Requested test streaming"); return mySink.asFlux() .startWith("INIT TEST") .doOnCancel(() -> log.info("On cancel")) .doOnTerminate(() -> log.info("On terminate")); } ... }
本地环境日志(正常触发)
2022-08-06 18:25:42.115 INFO 3685 --- [ main] com.wuase.sinkdemo.SinkDemoApplication : Starting SinkDemoApplication using Java 1.8.0_252 on aniello-pc with PID 3685 (/home/pc/eclipse-workspace/sink-demo/target/classes started by pc in /home/pc/eclipse-workspace/sink-demo) 2022-08-06 18:25:42.124 INFO 3685 --- [ main] com.wuase.sinkdemo.SinkDemoApplication : No active profile set, falling back to 1 default profile: "default" 2022-08-06 18:25:44.985 INFO 3685 --- [ main] o.s.b.web.embedded.netty.NettyWebServer : Netty started on port 8080 2022-08-06 18:25:45.018 INFO 3685 --- [ main] com.wuase.sinkdemo.SinkDemoApplication : Started SinkDemoApplication in 3.737 seconds (JVM running for 5.36) 2022-08-06 18:26:09.706 INFO 3685 --- [or-http-epoll-3] com.wuase.sinkdemo.TestController : Requested test streaming 2022-08-06 18:26:14.799 INFO 3685 --- [or-http-epoll-3] com.wuase.sinkdemo.TestController : On cancel
解决思路
调整Kubernetes Ingress/代理的超时配置:
大部分Ingress控制器(如Nginx Ingress)默认连接超时较长,客户端断开后不会即时通知后端。可以修改Ingress注解缩短超时,同时开启HTTP/1.1支持确保断开信号传递:apiVersion: networking.k8s.io/v1 kind: Ingress metadata: annotations: nginx.ingress.kubernetes.io/proxy-read-timeout: "30" nginx.ingress.kubernetes.io/proxy-send-timeout: "30" nginx.ingress.kubernetes.io/proxy-http-version: "1.1" nginx.ingress.kubernetes.io/proxy-set-header: "Connection \"\""配置Netty TCP保活参数:
让服务端主动检测死连接,避免TCP连接长时间处于半开状态。添加Netty自定义配置:@Bean public NettyServerCustomizer nettyServerCustomizer() { return server -> server.tcpConfiguration(tcp -> tcp.option(ChannelOption.SO_KEEPALIVE, true) .option(ChannelOption.TCP_KEEPIDLE, 300) // 5分钟无数据后启动检测 .option(ChannelOption.TCP_KEEPINTVL, 60) // 每隔1分钟发送探测包 .option(ChannelOption.TCP_KEEPCNT, 5) // 连续5次无响应则关闭连接 ); }改用doFinally覆盖所有终止场景:
doOnTerminate仅在流正常完成/错误终止时触发,doOnCancel仅处理取消场景,而doFinally会覆盖所有终止类型(取消、完成、错误),可以替换或补充原有逻辑:return mySink.asFlux() .startWith("INIT TEST") .doFinally(signal -> { if (signal == SignalType.CANCEL) { log.info("On cancel"); } else { log.info("On terminate: {}", signal); } });验证Sink的取消信号传递:
如果mySink是自定义多播Sink,确保开启自动取消配置,让订阅取消时及时传递信号:Sink<String> mySink = Sink.many() .multicast() .onBackpressureBuffer() .autoCancel(true);排查日志采集延迟/丢失:
有时不是方法未触发,而是K8s日志采集工具(如Loki、Fluentd)存在延迟。可以直接进入Pod查看本地日志文件:kubectl exec -it <pod-name> -- tail -f /path/to/app.log,确认事件是否真的未触发。检查NetworkPolicy规则:
确认K8s的NetworkPolicy没有阻止Pod与客户端之间的TCP FIN/RST包传递,确保相关规则允许双向TCP流量。
内容的提问来源于stack exchange,提问作者wuase

