Spring Integration动态WebSocket消息偶发未发送问题排查求助
环境与代码背景
基于Spring Boot 3.4.1 + Spring Integration 6.4.1开发,通过HTTP请求动态创建WebSocket,并从多线程向其发送消息。核心代码如下:
IntegrationFlowRegistration创建代码
private IntegrationFlowContext.IntegrationFlowRegistration buildWebSocketFlowRegistration(ServerWebSocketContainer serverWebSocketContainer, WebSocketAdapter webSocketAdapter) { WebSocketOutboundMessageHandler webSocketOutboundMessageHandler = new WebSocketOutboundMessageHandler(serverWebSocketContainer); StandardIntegrationFlow standardIntegrationFlow = IntegrationFlow.from(webSocketAdapter).split(new DecorateMessageWithSessionId(serverWebSocketContainer)).handle(webSocketOutboundMessageHandler).get(); IntegrationFlowContext.IntegrationFlowRegistration flowRegistration = integrationFlowContext.registration(standardIntegrationFlow).addBean(serverWebSocketContainer).register(); flowRegistration.start(); return flowRegistration; }
ServerWebSocketContainer创建代码
private ServerWebSocketContainer createWebSocketContainer(long channelID) { String[] allowedOrigins = Arrays.stream(applicationProperties.allowedOrigins.split(",")).map(String::trim).toArray(String[]::new); return new ServerWebSocketContainer("/" + channelID).setHandshakeHandler(handshakeHandler()).setAllowedOrigins(allowedOrigins); }
WebSocketAdapter实现
@Override public void sendMessage(Message<?> message) { super.sendMessage(message); }
异常现象
多数情况下消息发送正常,日志会完整显示从channel#0的preSend到后续splitter、outbound handler的处理流程;但偶发出现消息丢失:仅能看到channel#0的preSend日志,后续处理日志全无,消息未通过WebSocket传输。
已排查:WebSocket处于活跃状态,IntegrationFlowRegistration已启动,线程dump未发现阻塞情况,问题无固定复现规律。
可能原因与解决方案
1. 默认DirectChannel的消息丢弃特性
IntegrationFlow默认使用DirectChannel作为输入通道,它是同步阻塞型通道,当通道没有订阅者、或订阅者未就绪时,会直接丢弃消息(默认拦截器仅记录debug日志,不抛出异常)。多线程并发发送时,若此时流的订阅者(如splitter组件)未完全初始化,或订阅关系临时失效,就会出现仅preSend、无后续处理的情况。
解决方案:
- 将输入通道改为异步的
QueueChannel,它会缓存消息直到有订阅者处理:StandardIntegrationFlow standardIntegrationFlow = IntegrationFlow.from(webSocketAdapter, e -> e.channel(MessageChannels.queue())) .split(new DecorateMessageWithSessionId(serverWebSocketContainer)) .handle(webSocketOutboundMessageHandler) .get(); - 给DirectChannel添加自定义拦截器,当消息无法发送时抛出异常,便于定位:
DirectChannel inputChannel = new DirectChannel(); inputChannel.addInterceptor(new ChannelInterceptor() { @Override public Message<?> preSend(Message<?> message, MessageChannel channel) { return message; } @Override public void postSend(Message<?> message, MessageChannel channel, boolean sent) { if (!sent) { throw new MessageDeliveryException(message, "Failed to send message to channel: " + channel); } } }); StandardIntegrationFlow standardIntegrationFlow = IntegrationFlow.from(webSocketAdapter, e -> e.channel(inputChannel)) .split(new DecorateMessageWithSessionId(serverWebSocketContainer)) .handle(webSocketOutboundMessageHandler) .get();
2. WebSocketAdapter的线程安全隐患
WebSocketInboundChannelAdapter的父类实现可能存在线程不安全的状态,多线程并发调用sendMessage时,内部通道的订阅状态可能出现竞态,导致消息无法正确路由到后续流组件。
解决方案:
- 在
WebSocketAdapter的sendMessage方法添加同步锁,确保同一时间只有一个线程发送消息:@Override public synchronized void sendMessage(Message<?> message) { super.sendMessage(message); } - 确保每个WebSocket实例对应唯一的IntegrationFlow,避免多个流共享同一个适配器。
3. IntegrationFlow启动的异步初始化延迟
flowRegistration.start()是同步调用,但Spring Integration内部组件(如WebSocket会话注册、订阅者绑定)可能存在异步初始化逻辑。若在流完全启动前发送消息,消息会因无订阅者而被丢弃。
解决方案:
- 发送消息前验证流的运行状态,确保已完全启动:
if (flowRegistration.getIntegrationFlowContext().isRunning(flowRegistration.getId())) { webSocketAdapter.sendMessage(message); } else { // 处理流未就绪的情况,比如缓存消息或重试 log.warn("Integration flow not running, cannot send message"); }
4. Splitter组件的异常吞吃
自定义的DecorateMessageWithSessionId Splitter若在处理消息时抛出异常,且未配置错误处理逻辑,异常会被默认机制吞吃,导致消息中断,且可能仅在TRACE级别日志中记录。
解决方案:
- 给Splitter配置错误处理,捕获并记录异常:
StandardIntegrationFlow standardIntegrationFlow = IntegrationFlow.from(webSocketAdapter) .split(new DecorateMessageWithSessionId(serverWebSocketContainer)) .handle(webSocketOutboundMessageHandler) .errorHandler(throwable -> log.error("Error processing WebSocket message", throwable)) .get(); - 提高日志级别到
TRACE,查看Splitter组件的详细处理日志,确认是否有异常发生。
内容的提问来源于stack exchange,提问作者Andrei Roșu-Cojocaru

