Nifi ExecuteScript引用外部Jedis Jar运行不稳定问题排查求助
问题根因
- 静态资源生命周期不同步:自定义Redis类的静态
pool变量在onStop调用destroy()后没有被置空,下次调用getPoolInstance()时会直接返回已销毁的池,导致isClosed()返回true,这是功能时好时坏的核心原因。 - 类加载机制冲突:NiFi集群环境中,Groovy脚本的类加载器会随处理器启停、配置修改、节点状态同步重新实例化,自定义Jar的静态变量会绑定到不同的类加载器实例,出现多份池实例、状态不一致的问题。
- 连接池初始化逻辑缺陷:现有代码初始化JedisPool时未传入密码、超时等参数,且
getPoolConfig()用到的operationTimeout可能在未初始化的情况下被调用,导致参数异常。 - 全局引用冗余:Functions类中额外存储jedisPool静态引用,和Redis类的pool变量不同步,销毁操作仅修改了Functions的引用未同步到Redis类的静态变量。
修复方案
第一步:修改自定义Jar代码
1. 补充Redis类的销毁和安全校验逻辑
给静态变量加volatile修饰防止指令重排序,新增标准化销毁方法:
public class Redis { private static Object staticLock = new Object(); // 加volatile防止指令重排序导致半初始化实例被调用 private static volatile JedisPool pool; private static volatile JedisPoolConfig config; private static String host; private static int port; private static int connectTimeout; private static int operationTimeout; private static String password; // 新增销毁方法,销毁后自动清空静态变量 public static void destroyPool() { synchronized (staticLock) { if (pool != null) { pool.destroy(); pool = null; config = null; } } } // 初始化配置加锁校验,避免池初始化后被修改参数 public static void initializeSettings(String host, int port, String password, int connectTimeout, int operationTimeout) { synchronized (staticLock) { if (pool != null) { return; } Redis.host = host; Redis.port = port; Redis.password = password; Redis.connectTimeout = connectTimeout; Redis.operationTimeout = operationTimeout; } } public static JedisPool getPoolInstance() { if (pool == null) { synchronized(staticLock) { if (pool == null) { JedisPoolConfig poolConfig = getPoolConfig(); boolean useSsl = port == 6380; int db = 0; String clientName = "MyClientName"; SSLSocketFactory sslSocketFactory = null; SSLParameters sslParameters = null; HostnameVerifier hostnameVerifier = new SimpleHostNameVerifier(host); // 补全所有连接参数,避免密码、超时配置不生效 pool = new JedisPool(poolConfig, host, port, connectTimeout, operationTimeout, password, db, clientName, useSsl, sslSocketFactory, sslParameters, hostnameVerifier); } } } return pool; } public static JedisPoolConfig getPoolConfig() { if (config == null) { synchronized (staticLock) { if (config == null) { JedisPoolConfig poolConfig = new JedisPoolConfig(); int maxConnections = 200; poolConfig.setMaxTotal(maxConnections); poolConfig.setMaxIdle(maxConnections); poolConfig.setBlockWhenExhausted(true); poolConfig.setMaxWaitMillis(operationTimeout); poolConfig.setMinIdle(50); Redis.config = poolConfig; } } } return config; } // getPoolCurrentUsage、SimpleHostNameVerifier逻辑保持不变 }
2. 简化Functions类逻辑,移除冗余的JedisPool引用
public class Functions { private static final String UTF8= "UTF-8"; public static String searchPlace(double lattitude,double longitude) { // 直接从Redis类获取连接池,不额外存储静态引用避免状态不同步 try(Jedis jedis = Redis.getPoolInstance().getResource()) { // 编写你的GEO查询逻辑,比如调用jedis.georadius等命令 } catch(Exception e){ // 异常处理逻辑 } return ""; } }
第二步:修改Groovy脚本逻辑
import org.apache.nifi.processor.ProcessContext import com.customlib.functions.* import org.apache.commons.io.IOUtils import groovy.json.JsonSlurper import java.nio.charset.StandardCharsets def flowFile = session.get() if (flowFile == null) { return } def flowFiles = [] def failflowFiles = [] def input = null def data = null static onStart(ProcessContext context){ // 填写实际的Redis配置,超时时间不要传0,建议设置为3000~10000毫秒 Redis.initializeSettings("Redis服务地址", 6379, "Redis密码(无密码则填null)", 5000, 5000) } static onStop(ProcessContext context){ // 调用标准化销毁方法 Redis.destroyPool() } try { // 增加池状态校验,异常状态下自动重建连接池 if (Redis.getPoolInstance().isClosed()) { Redis.destroyPool() Redis.initializeSettings("Redis服务地址", 6379, "Redis密码(无密码则填null)", 5000, 5000) } log.warn('is jedispool connected::::' + !Redis.getPoolInstance().isClosed()) def inputStream = session.read(flowFile) def writer = new StringWriter() IOUtils.copy(inputStream, writer, "UTF-8") data = writer.toString() input = new JsonSlurper().parseText(data) def place = Functions.searchPlace(input["data"]["lat"] as double, input["data"]["longi"] as double) log.warn('place is::::' + place) // 业务逻辑处理 def newFlowFile = session.create(flowFile) newFlowFile = session.write(newFlowFile, { outputStream -> outputStream.write(data.getBytes(StandardCharsets.UTF_8)) } as OutputStreamCallback) flowFiles << newFlowFile } catch(Exception e) { log.error("处理异常", e) def newFlowFile = session.create(flowFile) newFlowFile = session.write(newFlowFile, { outputStream -> outputStream.write(data.getBytes(StandardCharsets.UTF_8)) } as OutputStreamCallback) failflowFiles << newFlowFile } finally { session.remove(flowFile) } session.transfer(flowFiles, REL_SUCCESS) session.transfer(failflowFiles, REL_FAILURE)
最优替代方案(推荐)
直接使用NiFi官方提供的RedisConnectionPoolService控制器服务,无需自行封装Jedis连接池:
- 在NiFi控制面板新增
RedisConnectionPoolService,配置Redis地址、端口、密码、超时等参数后启用服务 - 在Groovy脚本的处理器配置中添加对该控制器服务的引用
- 脚本中直接通过服务获取Jedis连接执行GEO查询即可,官方服务已经适配了集群环境、生命周期管理、多线程安全,稳定性远高于自定义实现。
内容的提问来源于stack exchange,提问作者ashok
相关产品推荐
相关产品推荐

