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

Flink集群模式集成Redisson异常重试时Redis连接数持续增长问题咨询

问题成因

  1. close方法执行异常与异步关闭不彻底:当前实现的close方法没有做非空校验,若open方法初始化Redisson过程中就抛出异常,redisson实例为null,调用shutdown时会触发空指针,导致close方法中断,旧连接资源无法释放。同时Redisson的shutdown()是异步方法,调用后不会立刻释放连接,若任务高频重启,新实例已经创建完成但旧连接还未完成回收,会出现连接数临时上涨的情况,长期累积就会超出上限。
  2. 集群模式与本地模式的运行差异:本地模式下所有任务组件运行在同一个JVM进程中,任务重试时线程、资源复用率极高,连接泄漏问题被掩盖。而集群模式下TaskManager的重启、算子subtask的重建都会触发open方法新建Redisson实例,若旧实例连接未被正确回收,连接数就会持续增长。
  3. Redisson 3.16.1版本存在资源泄漏bug:该版本的RX客户端在异常场景下调用shutdown时,存在连接池资源未完全释放的已知问题,进一步放大了连接泄漏的概率。

解决方案

1. 修复close方法的健壮性,保证资源同步释放

修改close方法逻辑,增加非空校验,调用同步等待方法确保连接完全回收:

@Override
public void close() throws Exception {
    // 增加非空校验避免空指针中断关闭逻辑
    if (redisson != null && !redisson.isShutdown()) {
        redisson.shutdown();
        // 最多等待3秒确保连接完全释放,可根据实际场景调整时长
        redisson.awaitShutdown(3, java.util.concurrent.TimeUnit.SECONDS);
        log.info("Shut down redisson instance successfully");
    }
    super.close();
}

2. 同TaskManager内复用Redisson实例

将Redisson实例改为静态持有,同一个TaskManager进程内的所有算子subtask复用同一个客户端,避免重复创建实例:

@Slf4j
public class ExceptionTestFunction extends KeyedProcessFunction<Integer, String, String> {
    // 静态实例,JVM级别唯一,同TM下所有subtask复用
    private static volatile RedissonRxClient redisson;
    // 计数引用,确保最后一个使用实例的subtask关闭时才释放资源
    private static volatile int refCount = 0;
    private static final Object LOCK = new Object();

    @Override
    public void open(Configuration parameters) {
        synchronized (LOCK) {
            if (redisson == null) {
                Config config = new Config();
                config.useSingleServer()
                        .setClientName("xhb-redisson-main")
                        .setTimeout(5000)
                        .setConnectTimeout(10000)
                        .setConnectionPoolSize(4)
                        .setConnectionMinimumIdleSize(2)
                        .setIdleConnectionTimeout(10000)
                        .setAddress("redis://127.0.0.1:6379")
                        .setDatabase(0)
                        .setPassword(null);
                redisson = Redisson.create(config).rxJava();
            }
            refCount++;
        }
    }

    @Override
    public void close() throws Exception {
        synchronized (LOCK) {
            refCount--;
            if (refCount == 0 && redisson != null && !redisson.isShutdown()) {
                redisson.shutdown();
                redisson.awaitShutdown(3, java.util.concurrent.TimeUnit.SECONDS);
                redisson = null;
                log.info("All subtasks exited, redisson instance shut down completely");
            }
        }
        super.close();
    }
}

3. 其他辅助优化

  • 升级Redisson版本到3.17.8及以上的稳定版,修复已知的资源泄漏问题;
  • 调整Flink重启策略,增加重启间隔,避免短时间内高频重启,给旧连接回收留足时间,示例配置:
    // 10分钟内最多重启3次,每次间隔10秒
    env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3, org.apache.flink.api.common.time.Time.seconds(10)));
    
  • 在Redis配置文件中调低tcp-keepalive参数到300秒,让Redis主动回收长时间无响应的僵死连接。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 15:36:02