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

Flink中CountMinSketch类ExecutorService序列化问题及解决方案咨询

解决Flink中ExecutorService不可序列化的问题

首先来看报错的核心原因:你的CountMinSketch类实现了Serializable接口,但它持有的ExecutorService实例(具体是Executors$FinalizableDelegatedExecutorService)本身不支持序列化。当Flink尝试序列化你的AverageAggregator(因为它持有CountMinSketch对象)时,就会触发这个NotSerializableException。

下面是针对这个问题的具体解决方案,同时还会处理异步更新带来的线程安全隐患:

方案一:延迟ExecutorService初始化到运行时(推荐)

Flink的富函数(比如RichAggregateFunction)提供了open()和close()方法,这两个方法会在任务实例启动后/销毁前调用,在这里创建的对象不需要被序列化。我们可以利用这个特性来管理ExecutorService的生命周期:

1. 修改CountMinSketch类

把ExecutorService标记为transient(避免序列化),并提供初始化方法:

import java.io.Serializable;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
import java.util.concurrent.atomic.AtomicInteger;

public class CountMinSketch implements Serializable {
    private static final long serialVersionUID = 1123747953291780413L;
    private static final int H1 = 0;
    private static final int H2 = 1;
    private static final int H3 = 2;
    private static final int H4 = 3;
    private static final int LIMIT = 100;
    // 改用AtomicInteger数组保证线程安全
    private final AtomicInteger[][] sketch = new AtomicInteger[4][LIMIT];
    final NaiveHashFunction h1 = new NaiveHashFunction(11, 9);
    final NaiveHashFunction h2 = new NaiveHashFunction(17, 15);
    final NaiveHashFunction h3 = new NaiveHashFunction(31, 65);
    final NaiveHashFunction h4 = new NaiveHashFunction(61, 101);
    // transient标记避免序列化
    private transient ExecutorService executor;

    public CountMinSketch() {
        // 初始化原子数组
        for (int i = 0; i < 4; i++) {
            for (int j = 0; j < LIMIT; j++) {
                sketch[i][j] = new AtomicInteger(0);
            }
        }
    }

    // 提供外部初始化ExecutorService的方法
    public void setExecutor(ExecutorService executor) {
        this.executor = executor;
    }

    public Future<Boolean> updateSketch(String value) {
        return executor.submit(() -> {
            // 原子操作更新,避免竞态条件
            sketch[H1][h1.getHashValue(value)].getAndIncrement();
            sketch[H2][h2.getHashValue(value)].getAndIncrement();
            sketch[H3][h3.getHashValue(value)].getAndIncrement();
            sketch[H4][h4.getHashValue(value)].getAndIncrement();
            return true;
        });
    }

    public Future<Boolean> updateSketch(String value, int count) {
        return executor.submit(() -> {
            sketch[H1][h1.getHashValue(value)].addAndGet(count);
            sketch[H2][h2.getHashValue(value)].addAndGet(count);
            sketch[H3][h3.getHashValue(value)].addAndGet(count);
            sketch[H4][h4.getHashValue(value)].addAndGet(count);
            return true;
        });
    }

    // 其他方法保持不变...
}

2. 修改AverageAggregator为富函数

继承RichAggregateFunction,在open()中创建ExecutorService,close()中销毁:

import org.apache.flink.api.common.functions.RichAggregateFunction;
import org.apache.flink.configuration.Configuration;
import java.util.concurrent.Executors;
import java.util.concurrent.ExecutorService;

public static class AverageAggregator extends RichAggregateFunction<Tuple3<Integer, Tuple5<Integer, String, Integer, String, Integer>, Double>, Tuple3<Double, Long, Integer>, Tuple2<String, Double>> {
    private static final long serialVersionUID = 7233937097358437044L;
    private String functionName;
    private CountMinSketch countMinSketch;
    // transient标记避免序列化
    private transient ExecutorService executor;

    @Override
    public void open(Configuration parameters) throws Exception {
        super.open(parameters);
        // 每个任务实例创建独立的ExecutorService
        executor = Executors.newSingleThreadExecutor();
        countMinSketch = new CountMinSketch();
        countMinSketch.setExecutor(executor);
    }

    @Override
    public void close() throws Exception {
        super.close();
        // 关闭ExecutorService,防止资源泄漏
        if (executor != null) {
            executor.shutdown();
        }
    }

    // 实现AggregateFunction的其他方法...
}

关键注意点

  1. 线程安全:原代码中的普通int数组在异步更新时会出现竞态条件,必须改用AtomicInteger数组或者加锁,上面的方案已经用原子类解决了这个问题。
  2. 资源管理:一定要在close()方法中关闭ExecutorService,避免分布式环境下的资源泄漏。
  3. Flink编程模型适配:自己管理线程池在Flink中需要谨慎,因为Flink本身有自己的资源调度机制,如果你的异步更新只是为了本地计算,上面的方案足够;如果涉及外部IO,建议使用Flink官方的Async I/O API,它更符合Flink的分布式特性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 04:29:52