提交Flink作业抛出JedisCluster序列化异常求助
解决Flink提交作业时JedisCluster的序列化异常问题
这个问题在Flink分布式场景里很常见,本质原因是:Flink需要把算子序列化后分发到各个TaskManager节点执行,但JedisCluster实例本身并没有实现Serializable接口。你的匿名KeyedProcessFunction引用了外部的jedisCluster对象,Flink尝试序列化这个算子时自然就抛出了NotSerializableException。
下面给你两种解决方案,优先推荐第一种,更贴合Flink的最佳实践:
方法1:在ProcessFunction生命周期方法中延迟初始化JedisCluster
不要在算子外部提前创建JedisCluster,而是把初始化逻辑放到算子的open方法里——这个方法是每个算子实例在TaskManager节点启动时执行的,完全不涉及序列化操作。同时记得在close方法中关闭资源,避免Redis连接泄漏。
修改后的代码如下:
DataStream<MobileClickEvent> clickEventDataStream = environment.addSource(clickConsumer); clickEventDataStream .filter(Objects::nonNull) .keyBy(new KeySelector<MobileClickEvent, String>() { @Override public String getKey(MobileClickEvent value) throws Exception { return value.getItemId() + "_" + value.getItemType(); } }) .process(new KeyedProcessFunction<String, MobileClickEvent, Object>() { // 用transient标记,告诉序列化框架跳过这个字段,不要尝试序列化它 private transient JedisCluster jedisCluster; @Override public void open(Configuration parameters) throws Exception { super.open(parameters); // 到了TaskManager节点上再初始化JedisCluster jedisCluster = JedisClusterBuilder.getInstance(JedisClusterEnum.THIRD); } @Override public void processElement(MobileClickEvent value, Context ctx, Collector<Object> out) throws Exception { String key = ctx.getCurrentKey(); jedisCluster.hincrBy("{item_feature}" + key, "click", 1); jedisCluster.expire("{item_feature}" + key, 60 * 10); } @Override public void close() throws Exception { super.close(); // 算子销毁前关闭Redis连接,防止资源泄漏 if (jedisCluster != null) { jedisCluster.close(); } } });
这里的核心要点:
transient关键字确保序列化框架不会处理jedisCluster字段open方法是算子在目标节点的初始化入口,在这里创建资源不会触发序列化问题close方法负责清理资源,符合Flink的算子生命周期管理规范
方法2:自定义序列化包装器(仅特殊场景使用)
如果因为业务限制必须在外部创建JedisCluster,可以自定义一个实现Serializable的包装类,序列化时只保存Redis集群的配置信息,反序列化时重新创建JedisCluster实例:
public class SerializableJedisCluster implements Serializable { private transient JedisCluster jedisCluster; private final JedisClusterEnum clusterEnum; public SerializableJedisCluster(JedisClusterEnum clusterEnum) { this.clusterEnum = clusterEnum; this.jedisCluster = JedisClusterBuilder.getInstance(clusterEnum); } // 序列化时只保存配置枚举,不处理JedisCluster实例 private void writeObject(ObjectOutputStream out) throws IOException { out.defaultWriteObject(); } // 反序列化时根据配置重新创建JedisCluster private void readObject(ObjectInputStream in) throws IOException, ClassNotFoundException { in.defaultReadObject(); this.jedisCluster = JedisClusterBuilder.getInstance(clusterEnum); } public JedisCluster getJedisCluster() { return jedisCluster; } }
然后在你的业务代码中使用这个包装类:
SerializableJedisCluster serializableJedisCluster = new SerializableJedisCluster(JedisClusterEnum.THIRD); // 后续process方法中通过包装类获取JedisCluster实例 .process(new KeyedProcessFunction<String, MobileClickEvent, Object>() { @Override public void processElement(MobileClickEvent value, Context ctx, Collector<Object> out) throws Exception { String key = ctx.getCurrentKey(); serializableJedisCluster.getJedisCluster().hincrBy("{item_feature}" + key, "click", 1); serializableJedisCluster.getJedisCluster().expire("{item_feature}" + key, 60 * 10); } });
不过这种方式不如第一种简洁,还可能因为配置不一致引发问题,所以优先推荐方法1。
内容的提问来源于stack exchange,提问作者xingbin
相关产品推荐
相关产品推荐

