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

Spring Integration动态TCP客户端IntegrationFlow未回收致堆内存溢出

问题分析与修复方案

核心问题

动态创建带UUID唯一标识的Spring Integration TCP客户端流后,调用IntegrationFlowContext.remove()无法彻底回收相关对象,导致堆内存溢出;业务要求每个请求对应全新TCP连接,因此必须使用唯一流ID。

直接可见的代码错误

你的TestController中调用removeFlow()时未传入对应的IntegrationTcpConnection对象,导致方法无法定位并移除目标流:

// 错误写法
manager.removeFlow();
// 正确写法
manager.removeFlow(con);

根本原因

  1. 未主动释放TCP底层资源:TcpNetClientConnectionFactory持有底层TCP连接,未调用stop()时,连接资源不会释放,导致上下文持有对象引用
  2. 自定义对象持有Spring组件引用:IntegrationTcpConnection持有MessageChannel、TcpNetClientConnectionFactory的强引用,即使流被移除,这些引用也会阻止GC回收
  3. 手动创建未被Spring管理的通道:手动创建的DirectChannel、QueueChannel未纳入Spring的生命周期管理,移除流后无法被自动清理

具体修复步骤

1. 修复流移除的调用错误

确保在测试代码中传入要移除的连接对象:

@GetMapping("/test")
public String test()
{
    for (int i = 0; i < 100; i++)
    {
        IntegrationTcpConnection con = manager.createConnection("localhost", 30000);
        con.send("test message");
        manager.removeFlow(con); // 传入连接对象
    }
    return "Success";
}

2. 主动关闭TCP连接工厂资源

在移除流前,先关闭连接工厂,释放底层TCP连接和资源:

public void removeFlow(IntegrationTcpConnection connection)
{
    // 先关闭连接工厂,释放底层TCP连接
    connection.getCf().stop();
    // 移除入站和出站流
    context.remove(connection.getBaseId() + ".in");
    context.remove(connection.getBaseId() + ".out");
    // 清空连接对象中的引用,帮助GC回收
    connection.clearReferences();
}

修改IntegrationTcpConnection类,添加清空引用的方法:

@AllArgsConstructor
public class IntegrationTcpConnection {

    private String baseId;
    private volatile MessageChannel send;
    private volatile PollableChannel receive;
    private volatile TcpNetClientConnectionFactory cf;

    public void send(String msg)
    {
        send.send(MessageBuilder.withPayload(msg).build());
    }

    public String receive(long timeout)
    {
        Message<?> message = receive.receive(timeout);
        return message != null ? (String) message.getPayload() : null;
    }

    public String getBaseId() {
        return baseId;
    }

    public TcpNetClientConnectionFactory getCf() {
        return cf;
    }

    // 清空引用,帮助GC回收
    public void clearReferences() {
        this.send = null;
        this.receive = null;
        this.cf = null;
    }
}

3. 改用Spring管理的通道,避免手动创建

不要手动创建DirectChannel,改用流注册后返回的输入通道,确保通道纳入Spring生命周期管理:

public IntegrationTcpConnection createConnection(String host, int port) {
    String baseId = host + port + UUID.randomUUID();
    TcpNetClientConnectionFactory cf = new TcpNetClientConnectionFactory(host, port);
    
    // 出站流:使用Spring管理的输入通道
    IntegrationFlow outFlow = f -> f.handle(Tcp.outboundAdapter(cf));
    IntegrationFlowRegistration outRegistration = context.registration(outFlow)
            .addBean(cf)
            .id(baseId + ".out")
            .register();
    MessageChannel sendChannel = outRegistration.getInputChannel();
    
    // 入站流:使用QueueChannel作为接收通道
    QueueChannel receiveChannel = new QueueChannel();
    IntegrationFlow inFlow = IntegrationFlow.from(Tcp.inboundAdapter(cf))
            .channel(receiveChannel)
            .get();
    context.registration(inFlow)
            .id(baseId + ".in")
            .register();

    return new IntegrationTcpConnection(baseId, sendChannel, receiveChannel, cf);
}

4. 优化入站适配器的启动逻辑

入站适配器默认自动启动,若不需要持续接收,可以设置autoStartup=false,按需启动:

IntegrationFlow inFlow = IntegrationFlow.from(Tcp.inboundAdapter(cf)
                        .autoStartup(false)) // 禁用自动启动
                .channel(receiveChannel)
                .get();
context.registration(inFlow)
        .id(baseId + ".in")
        .register();

// 当需要接收数据时,启动入站适配器
context.getRegistry().getIntegrationFlow(baseId + ".in")
        .getIntegrationComponents()
        .stream()
        .filter(comp -> comp instanceof TcpReceivingChannelAdapter)
        .map(TcpReceivingChannelAdapter.class::cast)
        .findFirst()
        .ifPresent(TcpReceivingChannelAdapter::start);

验证方法

修复后,使用内存分析工具(如IntelliJ内存分析器)检查:

  • 调用removeFlow()后,IntegrationFlow、TcpNetClientConnectionFactory、EventDrivenConsumer等对象是否被标记为可回收
  • 通过netstat命令检查是否存在未关闭的TCP连接
  • 多次调用测试端点,观察堆内存占用是否稳定

内容的提问来源于stack exchange,提问作者tyhennig

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 23:24:52