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

Spring Integration:如何通过网关从TCP服务器发消息并获取回复?

问题描述

我正尝试创建一个兼具TCP客户端与服务器功能的微服务用于测试。该微服务的客户端会连接到远程微服务的服务器端,同时远程微服务的客户端也会连接到本微服务的服务器端,从而实现任意一端均可发送消息并获取回复,通信链路如下:

Client side of μs1 <-> Server side of μs2 <-> Client side of μs2 <-> Server side of μs1

我曾尝试为客户端侧配置基于TcpSendingMessageHandler和TcpReceivingChannelAdapter的收发集成流,但由于这些是单向组件,无法等待回复并将其发送到replyChannel头,因此无法获取对方返回的消息,相关代码如下:

@Bean
public IntegrationFlow incomingClient(final TcpReceivingChannelAdapter tcpReceivingChannelAdapter,
                                      TcpServerEndpoint tcpServerEndpoint) {

    return IntegrationFlows
            .from(tcpReceivingChannelAdapter)
            .handle(message -> { LOGGER.info("RECEIVING ON CLIENT: {}", tcpServerEndpoint.processMessage((byte[]) message.getPayload()));})
            .get();
}

@Bean
public IntegrationFlow outgoingClient(final MessageChannel outboundChannel, final TcpSendingMessageHandler tcpSendingClientMessageHandler) {

    return IntegrationFlows
            .from(outboundChannel)
            .handle(tcpSendingClientMessageHandler)
            .get();
}

据我了解,需要使用TcpInboundGateway和TcpOutboundGateway组件来处理所需的回复。请问如何实现两端建立连接后,从服务器端发送消息并获取回复?能否通过InboundGateway以服务器身份发送消息?我需要实现任意一端均可发起通信并发送消息的功能。


解决方案

核心思路

要实现双向发起通信并获取回复,需同时配置TCP服务器网关(TcpInboundGateway)和TCP客户端网关(TcpOutboundGateway),两者配合完成双向请求-响应模式:

  • TcpInboundGateway:作为服务器端监听端口接收连接,处理远程客户端请求并返回回复;同时可借助连接会话存储主动向已连接客户端发送消息。
  • TcpOutboundGateway:作为客户端主动发起连接到远程服务器,发送请求并等待回复。

具体实现步骤

1. 配置TCP连接工厂

先定义客户端和服务器的连接工厂,确保双方使用一致的消息序列化/反序列化规则(示例用换行符分隔消息):

@Bean
public TcpNetClientConnectionFactory clientConnectionFactory() {
    TcpNetClientConnectionFactory factory = new TcpNetClientConnectionFactory("remote-host", 8080);
    factory.setSerializer(new ByteArrayCrLfSerializer());
    factory.setDeserializer(new ByteArrayCrLfSerializer());
    return factory;
}

@Bean
public TcpNetServerConnectionFactory serverConnectionFactory() {
    TcpNetServerConnectionFactory factory = new TcpNetServerConnectionFactory(9090);
    factory.setSerializer(new ByteArrayCrLfSerializer());
    factory.setDeserializer(new ByteArrayCrLfSerializer());
    return factory;
}

2. 配置TCP服务器网关(TcpInboundGateway)

该网关负责监听端口、处理请求并回复,同时通过TcpConnectionRepository存储连接,方便后续主动发消息:

@Bean
public TcpInboundGateway tcpInboundGateway(TcpNetServerConnectionFactory serverConnectionFactory,
                                           TcpConnectionRepository connectionRepository) {
    TcpInboundGateway gateway = new TcpInboundGateway();
    gateway.setConnectionFactory(serverConnectionFactory);
    gateway.setRequestChannel(serverRequestChannel());
    gateway.setConnectionRepository(connectionRepository);
    return gateway;
}

@Bean
public MessageChannel serverRequestChannel() {
    return new DirectChannel();
}

// 处理服务器收到的请求并返回回复
@ServiceActivator(inputChannel = "serverRequestChannel")
public byte[] handleServerRequest(byte[] payload) {
    String request = new String(payload);
    System.out.println("服务器收到请求: " + request);
    return ("服务器回复: " + request).getBytes();
}

3. 配置TCP客户端网关(TcpOutboundGateway)

该网关负责主动向远程服务器发送请求并等待回复:

@Bean
public TcpOutboundGateway tcpOutboundGateway(TcpNetClientConnectionFactory clientConnectionFactory) {
    TcpOutboundGateway gateway = new TcpOutboundGateway();
    gateway.setConnectionFactory(clientConnectionFactory);
    gateway.setRequestChannel(clientRequestChannel());
    return gateway;
}

@Bean
public MessageChannel clientRequestChannel() {
    return new DirectChannel();
}

// 示例:主动发送消息到远程服务器并获取回复
public String sendMessageToRemote(String message) {
    Message<byte[]> request = MessageBuilder.withPayload(message.getBytes())
            .setReplyChannelName("clientReplyChannel")
            .build();
    Message<byte[]> reply = (Message<byte[]>) clientRequestChannel().sendAndReceive(request);
    return new String(reply.getPayload());
}

@Bean
public MessageChannel clientReplyChannel() {
    return new DirectChannel();
}

4. 从服务器端主动发送消息

借助TcpConnectionRepository获取已建立的连接,通过TcpConnection主动向客户端发消息:

@Autowired
private TcpConnectionRepository connectionRepository;

public void sendMessageFromServer(String message) {
    // 获取所有已连接会话,可根据连接ID筛选特定客户端
    Collection<TcpConnection> connections = connectionRepository.getAllConnections();
    for (TcpConnection connection : connections) {
        connection.send(MessageBuilder.withPayload(message.getBytes()).build());
    }
}

关键说明

  • TcpInboundGateway本身是被动接收请求的组件,但通过TcpConnectionRepository可主动操作连接,实现服务器端主动推送消息。
  • 双向通信需两端同时部署上述网关:μs1的TcpOutboundGateway连接μs2的TcpInboundGateway,μs2的TcpOutboundGateway连接μs1的TcpInboundGateway,这样任意一端都能通过客户端网关发起请求,或通过服务器网关接收并回复请求。
  • 若需要主动发送消息后获取回复,需在客户端额外配置TcpReceivingChannelAdapter处理服务器主动发来的消息,形成完整的双向请求-响应通道。

内容的提问来源于stack exchange,提问作者Fernando Ramos Castillo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 14:31:51