Flink集群模式集成Redisson异常重试时Redis连接数持续增长问题咨询
问题成因
- close方法执行异常与异步关闭不彻底:当前实现的close方法没有做非空校验,若open方法初始化Redisson过程中就抛出异常,redisson实例为null,调用shutdown时会触发空指针,导致close方法中断,旧连接资源无法释放。同时Redisson的
shutdown()是异步方法,调用后不会立刻释放连接,若任务高频重启,新实例已经创建完成但旧连接还未完成回收,会出现连接数临时上涨的情况,长期累积就会超出上限。 - 集群模式与本地模式的运行差异:本地模式下所有任务组件运行在同一个JVM进程中,任务重试时线程、资源复用率极高,连接泄漏问题被掩盖。而集群模式下TaskManager的重启、算子subtask的重建都会触发open方法新建Redisson实例,若旧实例连接未被正确回收,连接数就会持续增长。
- 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
相关产品推荐
相关产品推荐

