Spring Integration获取远程源响应遇阻:配置ChannelAdapter未收到回复求助
我用Spring Integration实现和远程源的请求-响应通信,Wireshark能看到请求发出去、远程源也返回了回复,但应用就是拿不到回复消息。一开始用outboundGateway配置,因为回复消息没有EOF分隔符,而默认outboundGateway用的是ByteArrayCrLfSerializer,直接抛出了TimeoutException。后来改成用ChannelAdapter并设置ByteArrayRawSerializer,以为能拿到回复,但还是不行。
我的配置代码如下:
@Bean public AbstractClientConnectionFactory clientConnectionFactory() { TcpNioClientConnectionFactory connectionFactory = new TcpNioClientConnectionFactory(host, port); connectionFactory.setUsingDirectBuffers(true); connectionFactory.setSerializer(new ByteArrayRawSerializer()); connectionFactory.setSerializer(new ByteArrayRawSerializer()); connectionFactory.setApplicationEventPublisher(eventPublisher); return new CachingClientConnectionFactory(connectionFactory, poolSize); } @Bean public MessageChannel requestChannel() { return new DirectChannel(); } @Bean @ServiceActivator(inputChannel = "requestChannel", outputChannel = "responseChannel", requiresReply = "true") public TcpSendingMessageHandler requestHandler() { TcpSendingMessageHandler handler = new TcpSendingMessageHandler(); handler.setConnectionFactory(clientConnectionFactory()); return handler; } @Bean public MessageChannel responseChannel() { return new DirectChannel(); } @Bean public IntegrationFlow responseFlow() { return IntegrationFlow.from(responseChannel()) .handle(message -> { log.info("Reply: {}", message.getPayload()); }) .get(); }
我本来以为在@ServiceActivator上指定outputChannel后,就能在这个通道拿到回复,但实际根本没效果。现在找到的示例大多是默认的客户端-服务器连接配置,但因为回复是通过客户端打开的端口返回的,而且服务器只在客户端请求时才回复,我觉得不需要配置服务器连接,想问问是不是我理解错了,或者漏了什么配置?
问题原因与解决方案
1. TcpSendingMessageHandler的本质是单向发送器
TcpSendingMessageHandler是单向的出站通道适配器,它只负责发送消息,不会自动监听回复并将其转发到outputChannel。你在@ServiceActivator上配置outputChannel是无效的,因为这个Handler本身不产生回复消息。
2. 正确的配置方式:添加入站通道适配器监听回复
要接收远程源的回复,你需要单独配置一个TCP入站通道适配器,绑定到客户端连接工厂,用来监听已建立连接上返回的消息。
步骤1:修正连接工厂配置
去掉重复的setSerializer调用,同时明确设置反序列化器(ByteArrayRawSerializer可同时做序列化和反序列化):
@Bean public AbstractClientConnectionFactory clientConnectionFactory() { TcpNioClientConnectionFactory connectionFactory = new TcpNioClientConnectionFactory(host, port); connectionFactory.setUsingDirectBuffers(true); ByteArrayRawSerializer rawSerializer = new ByteArrayRawSerializer(); connectionFactory.setSerializer(rawSerializer); connectionFactory.setDeserializer(rawSerializer); connectionFactory.setApplicationEventPublisher(eventPublisher); return new CachingClientConnectionFactory(connectionFactory, poolSize); }
步骤2:添加TCP入站通道适配器
这个适配器会监听客户端连接上的回复消息,并转发到responseChannel:
@Bean @InboundChannelAdapter(channel = "responseChannel") public TcpReceivingChannelAdapter responseReceiver() { TcpReceivingChannelAdapter adapter = new TcpReceivingChannelAdapter(); adapter.setConnectionFactory(clientConnectionFactory()); return adapter; }
步骤3:调整请求Handler的配置
去掉@ServiceActivator上无效的outputChannel和requiresReply参数:
@Bean @ServiceActivator(inputChannel = "requestChannel") public TcpSendingMessageHandler requestHandler() { TcpSendingMessageHandler handler = new TcpSendingMessageHandler(); handler.setConnectionFactory(clientConnectionFactory()); return handler; }
3. 额外注意点
ByteArrayRawSerializer会读取所有可用字节直到连接关闭,如果远程源回复后不关闭连接,你需要自定义反序列化器来识别回复的结束标记(比如特定字节序列、固定长度等),否则入站适配器会一直阻塞等待更多数据。- 短连接场景(请求发送后远程源回复并关闭连接)下,
ByteArrayRawSerializer可以正常工作;长连接场景必须有明确的消息边界识别逻辑。
内容的提问来源于stack exchange,提问作者I have 10 fingers

