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

Spring Integration:闲置350秒后TCP连接断开的问题排查

问题描述

我们在AWS Kubernetes集群的多个Pod上运行Java Spring Integration应用,使用TCP Outbound Gateway与第三方系统通信,并通过CachingClientConnectionFactory缓存连接。已设置soKeepAlive为true,但闲置350秒后连接仍断开。是否需要额外配置,在闲置350秒前向服务器发送心跳以维持连接?

AWS相关限制:AWS VPC的NAT网关会主动断开闲置时间超过350秒的连接。

当前配置代码

@Bean
public AbstractClientConnectionFactory primeClientConnectionFactory() {
    TcpNetClientConnectionFactory tcpNetClientConnectionFactory = new TcpNetClientConnectionFactory(host, port);

    tcpNetClientConnectionFactory.setDeserializer(new PrimeCustomStxHeaderLengthSerializer());
    tcpNetClientConnectionFactory.setSerializer(new PrimeCustomStxHeaderLengthSerializer());
    tcpNetClientConnectionFactory.setSingleUse(false);
    tcpNetClientConnectionFactory.setSoKeepAlive(true);

    return tcpNetClientConnectionFactory;
}

@Bean
public AbstractClientConnectionFactory primeTcpCachedClientConnectionFactory() {
    CachingClientConnectionFactory cachingConnFactory = new CachingClientConnectionFactory(primeClientConnectionFactory(), connectionPoolSize);
    //cachingConnFactory.setSingleUse(false);
    cachingConnFactory.setLeaveOpen(true);
    cachingConnFactory.setSoKeepAlive(true);
    return cachingConnFactory;
}

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

@Bean
public RequestHandlerRetryAdvice retryAdvice() {
    RequestHandlerRetryAdvice retryAdvice = new RequestHandlerRetryAdvice();
    RetryTemplate retryTemplate = new RetryTemplate();
    FixedBackOffPolicy fixedBackOffPolicy = new FixedBackOffPolicy();
    fixedBackOffPolicy.setBackOffPeriod(500);
    SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy();
    retryPolicy.setMaxAttempts(3);
    retryTemplate.setBackOffPolicy(fixedBackOffPolicy);
    retryTemplate.setRetryPolicy(retryPolicy);
    retryAdvice.setRetryTemplate(retryTemplate);
    return retryAdvice;
}

@Bean
@ServiceActivator(inputChannel = "primeOutboundChannel")
public MessageHandler primeOutbound(AbstractClientConnectionFactory primeTcpCachedClientConnectionFactory) {
    TcpOutboundGateway tcpOutboundGateway = new TcpOutboundGateway();
    List<Advice> list = new ArrayList<>();
    list.add(retryAdvice());
    tcpOutboundGateway.setAdviceChain(list);

    tcpOutboundGateway.setRemoteTimeout(timeOut);
    tcpOutboundGateway.setRequestTimeout(timeOut);
    tcpOutboundGateway.setSendTimeout(timeOut);
    tcpOutboundGateway.setConnectionFactory(primeTcpCachedClientConnectionFactory);
    return tcpOutboundGateway;
}
解决方案

需要额外配置应用层心跳来维持连接,核心原因如下:

  • TCP的SO_KEEPALIVE默认超时逻辑(Linux下通常是2小时无数据后触发首个心跳)远长于AWS NAT网关的350秒闲置限制,无法提前触发保活。
  • NAT网关不会将TCP层面的KeepAlive包视为有效流量重置闲置计时器,只有应用层的业务心跳或探针包能维持连接。

方案1:实现Spring Integration TCP连接拦截器发送应用心跳

自定义拦截器定期发送心跳包,确保连接在350秒内保持活跃:

public class HeartbeatInterceptor implements TcpConnectionInterceptor {
    private final String heartbeatPayload;
    private final long intervalSeconds;
    private ScheduledExecutorService scheduler;

    public HeartbeatInterceptor(String heartbeatPayload, long intervalSeconds) {
        this.heartbeatPayload = heartbeatPayload;
        this.intervalSeconds = intervalSeconds;
    }

    @Override
    public void afterConnect(TcpConnection connection) {
        this.scheduler = Executors.newSingleThreadScheduledExecutor(r -> {
            Thread t = new Thread(r);
            t.setDaemon(true);
            return t;
        });
        // 间隔设置为300秒,小于350秒阈值
        this.scheduler.scheduleAtFixedRate(() -> {
            if (connection.isOpen()) {
                try {
                    connection.send(MessageBuilder.withPayload(heartbeatPayload).build());
                } catch (Exception e) {
                    connection.close();
                }
            } else {
                scheduler.shutdown();
            }
        }, intervalSeconds, intervalSeconds, TimeUnit.SECONDS);
    }

    @Override
    public void afterDisconnect(TcpConnection connection) {
        if (scheduler != null) {
            scheduler.shutdownNow();
        }
    }

    // 委托其他拦截器方法
    @Override
    public void setNextInterceptor(TcpConnectionInterceptor next) {}
    @Override
    public TcpConnectionInterceptor getNextInterceptor() { return null; }
    @Override
    public void onMessage(Message<?> message) {}
    @Override
    public void send(Message<?> message) {}
}

在连接工厂中注册拦截器:

@Bean
public AbstractClientConnectionFactory primeClientConnectionFactory() {
    TcpNetClientConnectionFactory tcpNetClientConnectionFactory = new TcpNetClientConnectionFactory(host, port);
    // 原有配置...
    tcpNetClientConnectionFactory.setInterceptorFactoryChain(chain -> 
        chain.addInterceptor(new HeartbeatInterceptor("PING", 300)));
    return tcpNetClientConnectionFactory;
}

方案2:配置连接池闲置超时自动重建

调整CachingClientConnectionFactory的闲置超时,让连接池在350秒前自动回收闲置连接,下次请求时重建:

@Bean
public AbstractClientConnectionFactory primeTcpCachedClientConnectionFactory() {
    CachingClientConnectionFactory cachingConnFactory = new CachingClientConnectionFactory(primeClientConnectionFactory(), connectionPoolSize);
    cachingConnFactory.setLeaveOpen(true);
    cachingConnFactory.setSoKeepAlive(true);
    // 设置闲置超时为300秒(300000毫秒)
    cachingConnFactory.setIdleTimeout(300000);
    return cachingConnFactory;
}

该方案无需修改第三方协议,但会带来连接重建的开销,适合对连接稳定性要求适中的场景。

额外注意事项

  • 确认第三方系统支持应用层心跳,避免心跳包被判定为无效请求。
  • 若使用Kubernetes探针,确保探针流量也能重置NAT网关的闲置计时器,或避免探针干扰连接状态。
  • 开启Spring Integration TCP模块的DEBUG日志(org.springframework.integration.tcp),便于排查连接生命周期问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 11:48:19