[Дополнительно] Несколько методов сегментированной обработки коллекций List в условиях многопоточности

задняя часть
[Дополнительно] Несколько методов сегментированной обработки коллекций List в условиях многопоточности

В последние два месяца из-за работы и семейных дел я был очень занят. У меня не так много времени, чтобы усваивать питательные вещества, и у меня нет выхода. В последнее время я немного расслабился, и последующие [ Advanced Road] будет медленно возвращаться на правильный путь.

Скажу прямо, первый раз, когда я столкнулся с многопоточной обработкой одной и той же задачи, было использование многопоточности IO для загрузки файлов, и с тех пор я не обрабатывал эту задачу, пока коллега не спросил меня несколько дней назад, почему многопоточная обработка списка Существуют различные ошибки с коллекциями, и как использовать несколько потоков для обработки одной и той же коллекции списка.

Во-первых, почему возникает проблема, похожая на повторную обработку определенного модуля?

Все мы знаем, что в Java каждый поток имеет свою независимую рабочую память, и все операции потока над общими переменными должны выполняться в его собственной рабочей памяти, и их нельзя читать и записывать напрямую из основной памяти.

Если измененное содержимое потока 1 должно быть получено потоком 2, измененная общая переменная в рабочей памяти потока 1 должна быть сначала обновлена ​​до основной памяти, а затем обновленная общая переменная в основной памяти обновлена ​​до рабочая память 2.

image.png

В настоящее время мы обычно рассматриваем использование различных методов синхронизации в java. Прежде всего, поскольку коллекция списков должна обрабатываться эффективно, синхронизированный метод можно исключить, поэтому я подумал об использовании CompletionService для выполнения асинхронных задач.

Подробные пояснения вы найдете в этой статье:[Дополнительно] Расширение пула потоков и CompletionService работают с асинхронными задачами.

1. Служба завершения

Сначала настройте WeedThreadPool, как в предыдущей статье.

public class WeedThreadPool extends ThreadPoolExecutor {
    private final ThreadLocal<Long> startTime =new ThreadLocal<>();
    private final Logger log =Logger.getLogger("WeedThreadPool");
    //统计执行次数
    private final AtomicLong numTasks =new AtomicLong();
    //统计总执行时间
    private final AtomicLong totalTime =new AtomicLong();
    /**
     * 这里是实现线程池的构造方法,我随便选了一个,大家可以根据自己的需求找到合适的构造方法
     */
    public WeedThreadPool(int corePoolSize, int maximumPoolSize, long keepAliveTime, TimeUnit unit, BlockingQueue<Runnable> workQueue) {
        super(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue);
    }
}

Затем есть метод реализации пула потоков для обработки коллекции списков.

public class WeedExecutorServiceDemo {
    BlockingQueue<Runnable> taskQueue;
    final static WeedThreadPool weedThreadPool = new WeedThreadPool(3, 10, 1, TimeUnit.SECONDS, new ArrayBlockingQueue<Runnable>(100));
    // 开始时间

    public static void main(String[] args) throws InterruptedException, ExecutionException {
        //记录任务开始时间
        long start = System.currentTimeMillis();
        CompletionService<List<Integer>> cs = new ExecutorCompletionService<>(weedThreadPool);
        int tb=1;
        //生成集合
        List<List<Integer>> list1 =new ArrayList();
        for (int i = 0; i < 10; i++) {
            List<Integer> list =new ArrayList();
            //随机生成任务处理
            int hb=tb;
            tb =tb*2;
            int finalTb = tb;
            cs.submit(new Callable<List<Integer>>(){

                @Override
                public List<Integer> call() throws Exception {
                    for (int j = hb; j< finalTb; j++){
                        list.add(j);
                    }
                    System.out.println(Thread.currentThread().getName()+"["+list+"]");

                    return list;
                }
            });
        }
        //注意在处理完毕后结束任务
        weedThreadPool.shutdown();
        for (int i = 0; i < 10; i++) {
            Future<List<Integer>> future = cs.take();
            if (future != null) {
                list1.add(future.get());
                System.out.println(future.get());
            }
        }
        System.err.println("执行任务消耗了 :" + (System.currentTimeMillis() - start) + "毫秒");
        System.out.println("結果["+list1.size()+"]==="+list1);
    }
}

результат процесса:

image.png

Судя по результатам, это довольно хорошо. CompletionService может относительно быстро обрабатывать задачи в сегментах. Я также упоминал ранее, что дизайн разумного размера пула потоков может помочь повысить эффективность обработки задач. Общий метод настройки в Интернете обычно такой: из:

Оптимальное количество потоков = ((время ожидания потока + процессорное время потока) / процессорное время потока) * количество процессоров

получить

Оптимальное количество потоков = (отношение времени ожидания потока к процессорному времени потока + 1) * количество процессоров

2. Вилка присоединения к пулу

Конечно, помимо использования CompletionService, вы также можете использовать ForkJoinPool для разработки метода обработки.

image.png

И ForkJoinPool, и ThreadPoolExecutor наследуются от абстрактного класса AbstractExecutorService, поэтому это почти то же самое, что использование ThreadPoolExecutor. Основная идея состоит в том, чтобы разделить большую задачу на несколько небольших задач, а затем объединить несколько небольших задач в один результат.

Платформа ForkJoinPool выполняет задачи, инициализируя ForkJoinTask, и предоставляет следующие два подкласса:

  • RecursiveAction: для задач, которые не возвращают результатов.
  • RecursiveTask: Используется для задач, которые возвращают результаты.

В процессе нашей реализации мы можем использовать метод RecursiveTask для обработки коллекции списков по секциям.

public class RecursiveTaskDemo {

    private static final ExecutorService executor = new ThreadPoolExecutor(2, 3, 10, TimeUnit.SECONDS, new LinkedBlockingQueue(10));
    private static final int totalRow = 53000;
    private static final int splitRow = 10000;

    public static void main(String[] args) throws InterruptedException, ExecutionException {
        long start = System.currentTimeMillis();
        //先循环生成待待处理集合
        List<Integer> list = new ArrayList<>(totalRow);
        for (int i = 0; i < totalRow; i++) {
            list.add(i);
        }
        //计算出需要创建的任务数
        int loopNum = (int)Math.ceil((double)totalRow/splitRow);
        ForkJoinPool pool = new ForkJoinPool(loopNum);
        ForkJoinTask<List> submit = pool.submit(new MyTask(list, 0, list.size()));

        List<List<Integer>>list1=new ArrayList<>();
        list1.add(submit.get());
        System.err.println("执行任务消耗了 :" + (System.currentTimeMillis() - start) + "毫秒");
        System.out.println("結果["+list1.size()+"]==="+list1);
    }
    //继承RecursiveTask
    static class MyTask extends RecursiveTask<List> {
        private List<Integer> list;
        private int startRow;
        private int endRow;

        public MyTask(List<Integer> list, int startRow, int endRow) {
            this.list = list;
            this.startRow = startRow;
            this.endRow = endRow;
        }

        /**
         * 递归处理数据,计算
         * @return
         */
        @Override
        protected List compute() {
            if (endRow - startRow <= splitRow) {
                List<Integer> ret = new ArrayList<>();
                for (int i = startRow; i < endRow; i++) {
                    //递归处理数据
                    ret.add(list.get(i));
                }
                System.out.println(Thread.currentThread().getName()+"["+ret+"]");
                return ret;
            }
            int loopNum = (int)Math.ceil((double)totalRow/splitRow);
            int startRow = 0;
            List<MyTask> myTaskList = new ArrayList<>();
            for (int i = 0; i < loopNum; i++) {
                if (startRow > totalRow) {
                    break;
                }
                int endRow = Math.min(startRow + splitRow, totalRow);
                System.out.println(String.format("startRow:%s, endRow:%s", startRow, endRow));
                myTaskList.add(new MyTask(list, startRow, endRow));
                startRow += splitRow;
            }
            //调用不同线程上独立执行的任务
            invokeAll(myTaskList);
            List<Integer> ret = new ArrayList<>();
            //归并
            for (MyTask myTask : myTaskList) {
                ret.addAll(myTask.join());
            }
            return ret;
        }
    }
}

результат процесса:

image.png

С помощью метода, показанного выше, вы можете повысить эффективность обработки задач разблокированным способом.Его можно использовать в таких сценариях, как обработка данных сканером и миграция данных, и измеренный эффект неплох. Конечно, судя по анализу результатов обработки, эффективность CompletionService наверняка выше~.

Всем привет,я Нанджу практикую джаву два с половиной года.Следующее мой WeChat.Друзья,кому нужны предыдущие карты или хотят обменяться опытом,могут общаться друг с другом.