【问题标题】:ExecutorService of Runnable, processing Batches of ArrayList not finishing processingRunnable的ExecutorService,处理ArrayList的Batches未完成处理
【发布时间】:2020-04-20 01:09:12
【问题描述】:

编辑:谢谢马克,对于那些有类似问题的人,我的问题是我先创建了一个可运行类的 Thread 实例,然后将线程提交给 executorservice。
它帮助我弄清楚,实际上,当我使用 ExecutorService 时,是否有未捕获的异常;它不会通知您,它将取消该过程,没有通知。这就是我处理不完整的原因。

我有一个对象的 ArrayList,我想批量处理多线程,但限制在给定时间运行的线程数。我发现 ExecutorService 可以处理这个问题。但是在测试它是否正在处理每条记录时,它似乎只处理了我传递给它的对象的一小部分。

编辑:我已经删除了它的多线程部分,并在不使用执行器服务的情况下像平常一样处理小批量(仅 710)的对象,它工作正常;线程是否有可能完成得太快并且处理不正确?这意味着通常一次处理大约 300k-800k 条记录;这就是为什么我想多线程。

public void processContainerRecords(ArrayList<? extends ContainerRecord> records) {
    int cores = Runtime.getRuntime().availableProcessors();
    ExecutorService executor = Executors.newFixedThreadPool(cores);
    int batchSize = Settings.LOGIC_BATCH_SIZE;//100
    int batches = (int) Math.ceil((double) records.size() / (double) batchSize);

    ArrayList<Future<?>> threads = new ArrayList<Future<?>>();
    LogicProcessor newHandler = null;
    for (int startIndex = 0; startIndex < records.size(); startIndex += batchSize + 1) {
        if (records.size() < batchSize) {
            newHandler = new LogicProcessor(mainGUI, records.subList(startIndex, records.size()));
        } else {
            int bound = (startIndex + batchSize);
            if (bound > records.size()) {
                bound = records.size();
            }
            newHandler = new LogicProcessor(mainGUI, records.subList(startIndex, bound));
        }
        Thread newThread = new Thread(newHandler);
        Future<?> f = executor.submit(newThread);
        threads.add(f);
    }
    executor.shutdown();
    int completedThreads = 0;
    while (!executor.isTerminated()) {//monitors threads and waits until completion
        completedThreads = 0;
        for (Future<?> f : threads) {
            if (f.isDone()) {
                completedThreads++;
            }
        }
        //currentProgress = completedThreads;
    }

    for (ContainerRecord record : records) {//checks if each record has been processed
        System.out.println(record.getContainer() + ":" + record.isTouched());
    }
}

这是启动线程实例的 LogicProcessor 类

    private List<? extends ContainerRecord> archive;
private GUI mainGUI;

public LogicProcessor(GUI mainGUI, List<? extends ContainerRecord> records) {
    this.mainGUI = mainGUI;
    this.archive = records;
}

@Override
public void run() {
    handleLogic();
}

private void handleLogic() {
    Iterator iterator = archive.iterator();
    while (iterator.hasNext()) {
        ContainerRecord record = (ContainerRecord) iterator.next();
        record.touch();//sets a boolean in the object to validate if it has been processed yet.
    }
}

输出:在处理的 710 条记录(对象)中,691 条从未被处理/触摸过,只有 19 条已经处理。

这有什么问题?我已经尝试了很多方法,甚至制作了一个类 LogicProcessor 的数组并将实例保存在数组中以避免任何类型的 GC 删除实例。我不确定它为什么不处理这些记录。

【问题讨论】:

    标签: java arraylist batch-processing executorservice


    【解决方案1】:

    我现在没有计算机来运行测试,但是通过查看您的代码,所以我的回答是基于个人经验的,甚至可以看作是代码审查,因为缺乏代码清晰度是错误的根源 :)

    1. 不要将new Thread 提交到执行器服务中。执行器的全部目的是向用户隐藏带有线程的单词。相反,您的 LogicProcessor 应该实现 Runnable/Callable 接口,具体取决于您是否要返回值。

    2. 再次检查分批分区的逻辑。如果您使用番石榴,它已经实现了分区逻辑。见this tutorial。我承认这更多是个人喜好,你的代码可能也很好,我还没有深入检查。

    3. 可能会简化关闭方法和期货处理。

    调用shutdown 方法会导致执行器服务停止接受新的任务来执行,但它不会立即关闭服务,而是会等待它已经拥有的所有任务都被执行。通常这样的线程池是在应用程序生命周期的开始时创建的,并且只要应用程序运行就一直存在。创建一个池非常昂贵,因为它分配线程。

    如果您想保持池打开但确保所有任务都已完成,您可以使用循环来迭代未来。

    所以我真的没有理由同时使用两者。如果您分配池只是为了提交一堆任务 - 调用 shutdown 就足够了。否则,您可以使用循环并将池视为全局对象并在其他地方调用 shutdown,正如我在上面解释的那样。

    【讨论】:

    • 太棒了,只是跳过让它成为一个线程,只是将实现 Runnable 的类传递到执行器服务中,使它完美地处理了所有事情!谢谢!
    猜你喜欢
    • 2012-11-28
    • 2019-01-05
    • 1970-01-01
    • 2019-01-13
    • 2018-05-03
    • 2021-03-18
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多