Создайте простой и удобный в использовании параллельный компонент на основе JDK ForkJoin.

задняя часть

Создайте простой и удобный в использовании параллельный компонент на основе ForkJoin.

В реальном развитии бизнеса необходимо использовать знания параллельного программирования, фактическое использование пула асинхронных потоков для выполнения задач сцены не особенно велико, и, как правило, действительно существует потребность в параллельном использовании, может быть более распространенным является прямое реализация интерфейса Runnable/Callable для выполнения броскового потока; или немного более продвинутая, определение пула потоков, брошенных для выполнения; фильм Боуэна, с другого ракурса, с ForkJoin JDK, предоставленным для разработки простой в использовании структуры для параллелизма

I. Предыстория

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

// 获取商品基本信息
ItemInfo itemInfo = itemService.getInfo(itemId);

// 获取销量
int sellCount = sellService.getSellCount(itemId);

// 获取评价信息
RateInfo rateInfo = rateService.getRateInfo(itemId);


// 获取店铺信息
ShopInfo shopInfo = shopService.getShopInfo(shopId);


// 获取装饰信息
DecorateInfo decoreateInfo = decorateService.getDecorateInfo(itemId);

// 获取推荐商品
RecommandInfo recommandInfo = recommandService.getRecommand(itemId);

Если это обычный процесс выполнения, то вышеупомянутые 6 вызовов выполняются последовательно.Предполагая, что rt каждой службы составляет 10 мс, тогда время выполнения шести служб здесь> 60 мс.

Но с точки зрения бизнеса вышеупомянутые шесть вызовов службы никак не связаны друг с другом, то есть вызов одной службы не зависит от результата, возвращаемого другой службой, и они могут выполняться одновременно. , выполнение шести сервисов требует времени, это тот, который занимает больше всего времени среди шести сервисов, возможно, чуть больше 10 мс.

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

II. Разработка и реализация

Например, в приведенном выше случае, если мы используем пул резьбы и как вы можете этого добиться?

1. Метод пула потоков

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

// 1. 创建线程池
ExecutorService alarmExecutorService = new ThreadPoolExecutor(3, 5, 60,
                TimeUnit.SECONDS,
                new LinkedBlockingDeque<>(10), 
                new DefaultThreadFactory("service-pool"),
                new ThreadPoolExecutor.CallerRunsPolicy());


// 2. 将服务调用,封装到线程任务中执行
Future<ItemInfo> itemFuture = alarmExecutorService.submit(new Callable<ItemInfo>() {
    @Override
    public ItemInfo call() throws Exception {
        return itemService.getInfo(itemId);
    }
});

// ... 其他的服务依次类推


// 3. 获取数据
ItemInfo = itemFutre.get(); // 阻塞,直到返回

Можно сказать, что приведенная выше реализация является очень четкой реализацией.Давайте посмотрим, как играть с инфраструктурой Fork/Join и каковы ее преимущества.

2. Метод ForkJoin

Прежде всего, вам, возможно, потребуется кратко представить, что это такое.Среда Fork/Join — это структура, предоставляемая Java7 для параллельного выполнения задач.Он делит большую задачу на несколько маленьких задач и, наконец, суммирует результаты каждой маленькой задачи. Фреймворк для получения результатов больших задач после

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

Заметки об использовании обучения ForkJoin

Как с помощью ForkJoin мы можем поддерживать описанные выше сценарии? Простое решение заключается в следующем

// 1. 创建池
ForkJoinPool pool = new ForkJoinPool(10);


// 2. 创建任务并提交
ForkJoinTask<ItemInfo> future = joinPool.submit(new RecursiveTask<ItemInfo>() {
    public ItemInfo compute() {
        return itemService.getItemInfo(itemId);
    }
});


// 3. 获取结果
future.join();

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

3. Расширенный

Как в полной мере использовать идею дизассемблирования задач ForkJoin для решения проблемы?

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

arch

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

Таким образом, текущая цель состоит в том, как реализовать приведенные выше требования сценария разделения задач, а также надеяться, что существующие изменения кода не слишком велики, ключ в том, что после его написания его должно быть легко понять + поддерживать (это на самом деле очень важно, автор столкнулся с ситуацией, когда упаковка особенно хороша, что приводит к слишком большим затратам на обслуживание при передаче бизнеса и резкому усложнению устранения неполадок)

4. Реализация

А. Идеи дизайна

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

Из-за разборки задачи нам нужна специальная задача, которая может быть набором из нескольких задач (то есть большая задача, сначала называемая bigTask)

Затем, когда он используется, все задачи инкапсулируются в bigTask и напрямую перебрасываются в forkJoinPool для выполнения (поддерживает метод вызова вызова для синхронного получения результата и метод выполнения для асинхронного получения результата).

Таким образом, суть заключается в том, как спроектировать эту Большую Задачу, и при ее выполнении ядро ​​разбирается на более детализированную Большую Задачу или Задачу и, в конечном итоге, объединяет все результаты выполнения ЗАДАЧИ и возвращает результат.

б. осознать

Основной интерфейс задачи

/**
 * Created by yihui on 2018/4/8.
 */
public interface IDataLoader<T> {


    /**
     * 具体的业务逻辑,放在这个方法里面执行,将返回的结果,封装到context内
     *
     * @param context
     */
    void load(T context);

}

Класс абстрактной реализации, наследующий ReuriAction от forkjoin, который соответствует базовой задаче, которую мы определили ранее.

public abstract class AbstractDataLoader<T> extends RecursiveAction implements IDataLoader {

    // 这里就是用来保存返回的结果,由业务防自己在实现的load()方法中写入数据
    protected T context;

    public AbstractDataLoader(T context) {
        this.context = context;
    }

    public void compute() {
        load(context);
    }


    /**
     * 获取执行后的结果,强制等待执行完毕
     * @return
     */
    public T getContext() {
        this.join();
        return context;
    }

    public void setContext(T context) {
        this.context = context;
    }
}

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

public class DefaultForkJoinDataLoader<T> extends AbstractDataLoader<T> {
    /**
     * 待执行的任务列表
     */
    private List<AbstractDataLoader> taskList;


    public DefaultForkJoinDataLoader(T context) {
        super(context);
        taskList = new ArrayList<>();
    }


    public DefaultForkJoinDataLoader<T> addTask(IDataLoader dataLoader) {
        taskList.add(new AbstractDataLoader(this.context) {
            @Override
            public void load(Object context) {
                dataLoader.load(context);
            }
        });
        return this;
    }


    // 注意这里,借助fork对任务进行了拆解
    @Override
    public void load(Object context) {
        this.taskList.forEach(ForkJoinTask::fork);
    }


    /**
     * 获取执行后的结果
     * @return
     */
    public T getContext() {
        this.taskList.forEach(ForkJoinTask::join);
        return this.context;
    }
}

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

public class ExtendForkJoinPool extends ForkJoinPool {

    public ExtendForkJoinPool() {
    }

    public ExtendForkJoinPool(int parallelism) {
        super(parallelism);
    }

    public ExtendForkJoinPool(int parallelism, ForkJoinWorkerThreadFactory factory, Thread.UncaughtExceptionHandler handler, boolean asyncMode) {
        super(parallelism, factory, handler, asyncMode);
    }


    // 同步阻塞调用时,需要对每个task执行join,确保执行完毕
    public <T> T invoke(ForkJoinTask<T> task) {
        if (task instanceof AbstractDataLoader) {
            super.invoke(task);
            return (T) ((AbstractDataLoader) task).getContext();
        } else {
            return super.invoke(task);
        }
    }
}

Затем есть фабричный класс, который создает пул, ничего особенного.

public class ForkJoinPoolFactory {

    private int parallelism;

    private ExtendForkJoinPool forkJoinPool;

    public ForkJoinPoolFactory() {
        this(Runtime.getRuntime().availableProcessors() * 16);
    }

    public ForkJoinPoolFactory(int parallelism) {
        this.parallelism = parallelism;
        forkJoinPool = new ExtendForkJoinPool(parallelism);
    }

    public ExtendForkJoinPool getObject() {
        return this.forkJoinPool;
    }

    public int getParallelism() {
        return parallelism;
    }

    public void setParallelism(int parallelism) {
        this.parallelism = parallelism;
    }


    public void destroy() throws Exception {
        this.forkJoinPool.shutdown();
    }

}

На этом все в основном закончилось, каждый класс очень прост, всего несколько вещей, следующий шаг — посмотреть, как его использовать.

III. Тестовая проверка

Давайте начнем с простого случая, чтобы продемонстрировать, как его использовать.

 @Data
static class Context {
    public int addAns;

    public int mulAns;

    public String concatAns;

    public Map<String, Object> ans = new ConcurrentHashMap<>();
}


@Test
public void testForkJoinFramework() {
    ForkJoinPool forkJoinPool = new ForkJoinPoolFactory().getObject();

    Context context = new Context();
    DefaultForkJoinDataLoader<Context> loader = new DefaultForkJoinDataLoader<>(context);
    loader.addTask(new IDataLoader<Context>() {
        @Override
        public void load(Context context) {
            context.addAns = 100;
            System.out.println("add thread: " + Thread.currentThread());
        }
    });
    loader.addTask(new IDataLoader<Context>() {
        @Override
        public void load(Context context) {
            try {
                Thread.sleep(3000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            context.mulAns = 50;
            System.out.println("mul thread: " + Thread.currentThread());
        }
    });
    loader.addTask(new IDataLoader<Context>() {
        @Override
        public void load(Context context) {
            context.concatAns = "hell world";
            System.out.println("concat thread: " + Thread.currentThread());
        }
    });


    DefaultForkJoinDataLoader<Context> subTask = new DefaultForkJoinDataLoader<>(context);
    subTask.addTask(new IDataLoader<Context>() {
        @Override
        public void load(Context context) {
            System.out.println("sub thread1: " + Thread.currentThread() + " | now: " + System.currentTimeMillis());
            try {
                Thread.sleep(200);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            context.ans.put(Thread.currentThread().getName(), System.currentTimeMillis());

        }
    });
    subTask.addTask(new IDataLoader<Context>() {
        @Override
        public void load(Context context) {
            System.out.println("sub thread2: " + Thread.currentThread() + " | now: " + System.currentTimeMillis());
            context.ans.put(Thread.currentThread().getName(), System.currentTimeMillis());
        }
    });

    loader.addTask(subTask);


    long start = System.currentTimeMillis();
    System.out.println("------- start: " + start);

    // 提交任务,同步阻塞调用方式
    forkJoinPool.invoke(loader);


    System.out.println("------- end: " + (System.currentTimeMillis() - start));

    // 输出返回结果,要求3s后输出,所有的结果都设置完毕
    System.out.println("the ans: " + context);
}

Он относительно прост в использовании, всего четыре простых шага:

  • Создать пул
  • Указывает класс контейнера ContextHolder для сохранения результата
  • Создать задачу
    • Создать корневую задачуnew DefaultForkJoinDataLoader<>(context);
    • Добавить подзадачу
  • Отправить

В приведенной выше реализации будет очень просто снова разделить задачу, посмотрите на приведенный выше вывод.

------- start: 1523200221827
add thread: Thread[ForkJoinPool-1-worker-50,5,main]
concat thread: Thread[ForkJoinPool-1-worker-36,5,main]
sub thread2: Thread[ForkJoinPool-1-worker-29,5,main] | now: 1523200222000
sub thread1: Thread[ForkJoinPool-1-worker-36,5,main] | now: 1523200222000
mul thread: Thread[ForkJoinPool-1-worker-43,5,main]
------- end: 3176
the ans: ForJoinTest.Context(addAns=100, mulAns=50, concatAns=hell world, ans={ForkJoinPool-1-worker-36=1523200222204, ForkJoinPool-1-worker-29=1523200222000})
  • Первый - это вывод потока, выполняемого каждой подзадачей. Видно, что это действительно задача, выполняемая разными потоками (параллелизм).
  • Через 3 с результат вывода, то есть после вызова, будет заблокирован до тех пор, пока не будут выполнены все задачи.
  • Подзадада разбирается на задачу. Время выполнения двух подзадач совпадает, но один сон, а другой не затронут (подзадачи также выполняются параллельно)

Для случая, когда вы хотите выполнить асинхронно, это тоже относительно просто, достаточно изменить место отправки задачи, а затем, когда вам нужно будет получить данные, вы сможете получить результат через загрузчик.

@Test
public void testForkJoinFramework2() {
    ForkJoinPool forkJoinPool = new ForkJoinPoolFactory().getObject();

    Context context = new Context();
    DefaultForkJoinDataLoader<Context> loader = new DefaultForkJoinDataLoader<>(context);
    loader.addTask(new IDataLoader<Context>() {
        @Override
        public void load(Context context) {
            try {
                Thread.sleep(3000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            context.addAns = 100;
            System.out.println("add thread: " + Thread.currentThread());
        }
    });
    loader.addTask(new IDataLoader<Context>() {
        @Override
        public void load(Context context) {
            context.mulAns = 50;
            System.out.println("mul thread: " + Thread.currentThread());
        }
    });
    loader.addTask(new IDataLoader<Context>() {
        @Override
        public void load(Context context) {
            context.concatAns = "hell world";
            System.out.println("concat thread: " + Thread.currentThread());
        }
    });


    long start = System.currentTimeMillis();
    System.out.println("------- start: " + start);

    // 如果暂时不关心返回结果,可以采用execute方式,异步执行
    forkJoinPool.execute(loader);

    // .... 这里可以做其他的事情 此时,不会阻塞,addAns不会被设置
    System.out.println("context is: " + context);
    System.out.println("------- then: " + (System.currentTimeMillis() - start));


    loader.getContext(); // 主动调用这个,表示会等待所有任务执行完毕后,才继续下去
    System.out.println("context is: " + context);
    System.out.println("------- end: " + (System.currentTimeMillis() - start));
}

IV. Другое

исходный код

Соответствующий исходный код можно посмотреть на git, в основном в проекте Quick-Alarm.

личный блог:Серый блог

Личный блог, записывайте все посты в блоге по учебе и работе, приглашаю всех в гости

утверждение

Это не так хорошо, как письмо веры.Контент был опубликован, и это чисто из семьи.Из-за моих средних способностей и ограниченных знаний, если вы обнаружите ошибки или у вас есть лучшие предложения, вы всегда можете критиковать и поправьте их.

Сканировать внимание

QrCode