Фабричный класс, предоставляющий Executor
Игнорировать пользовательские ThreadFactory, вызываемые и ненастраиваемые связанные методы
-
newFixedxxx: В любой момент существует не более nThreads потоков, обрабатывающих задачу; если новая задача поступает, когда все потоки работают, она будет помещена в очередь; если поток по какой-то причине завершается во время выполнения, если последующее действие задача должна быть выполнена, новый поток заменит ее
return new ThreadPoolExecutor(nThreads, nThreads, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<Runnable>()); -
newCachedxxx: при поступлении новой задачи она будет повторно использована, если в пуле потоков есть незанятые потоки, в противном случае будет создан новый поток. Если поток не используется более 60 секунд, он будет закрыт и удален из пула потоков.
return new ThreadPoolExecutor(0, Integer.MAX_VALUE, 60L, TimeUnit.SECONDS, new SynchronousQueue<Runnable>()); -
newSingleThreadExecutor: Для обработки задач используется только один поток, если этот поток зависнет, для его замены будет сгенерирован новый поток. Каждое задание гарантированно выполняется по порядку и только по одному
public static ExecutorService newSingleThreadExecutor() { return new FinalizableDelegatedExecutorService (new ThreadPoolExecutor(1, 1, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<Runnable>())); }Использование метода newFixedxxx также может привести к аналогичному эффекту, но ThreadPoolExecutor предоставит метод для изменения количества потоков, а FinalizableDelegatedExecutorService не имеет возможности его изменить, он основан на DelegatedExecutorService. Вышеуказанное предоставляет только метод для закрытия потока при выполнении finalize, а DelegatedExecutorService предоставляет только метод самого ExecutorService.
-
newScheduledThreadPool: предоставляет пул потоков для задержки или периодического выполнения задач.
public ScheduledThreadPoolExecutor(int corePoolSize) { super(corePoolSize, Integer.MAX_VALUE, 0, TimeUnit.NANOSECONDS, new DelayedWorkQueue()); } -
newSingleThreadScheduledExecutor: предоставляет один поток для задержки или периодического выполнения задач.Если исполняемый поток зависает, будет создан новый.
return new DelegatedScheduledExecutorService (new ScheduledThreadPoolExecutor(1));Кроме того, это гарантирует, что количество потоков самого возвращаемого Executor не может быть изменено.
Как видно из приведенной выше реализации, ядро состоит из трех частей.
- ThreadPoolExecutor: обеспечивает управление, связанное с количеством потоков.
- DelegatedExecutorService: предоставляет только методы самого ExecutorService, гарантируя, что количество потоков останется неизменным для достижения семантических сценариев.
- ScheduledExecutorService: Предоставляет функцию отложенного или периодического выполнения
Соответственно, есть и разные очереди для реализации разных сценариев.
- LinkedBlockingQueue: неограниченная очередь блокировки
- SynchronousQueue: при отсутствии потребления потребителем новые задачи будут заблокированы.
- DelayQueue: задачи в очереди могут быть выполнены после истечения срока их действия, в противном случае элементы в очереди не могут быть запрошены.
DelegatedExecutorService
Он просто обертывает метод ExecutorService и передает его входящему ExecutorService для выполнения.Так называемый UnConfigurable на самом деле означает, что он не предоставляет метод настройки различных настроек параметров.
static class DelegatedExecutorService extends AbstractExecutorService {
private final ExecutorService e;
DelegatedExecutorService(ExecutorService executor) { e = executor; }
public void execute(Runnable command) { e.execute(command); }
public void shutdown() { e.shutdown(); }
public List<Runnable> shutdownNow() { return e.shutdownNow(); }
public boolean isShutdown() { return e.isShutdown(); }
public boolean isTerminated() { return e.isTerminated(); }
public boolean awaitTermination(long timeout, TimeUnit unit)
throws InterruptedException {
return e.awaitTermination(timeout, unit);
}
public Future<?> submit(Runnable task) {
return e.submit(task);
}
public <T> Future<T> submit(Callable<T> task) {
return e.submit(task);
}
public <T> Future<T> submit(Runnable task, T result) {
return e.submit(task, result);
}
public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks)
throws InterruptedException {
return e.invokeAll(tasks);
}
public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks,
long timeout, TimeUnit unit)
throws InterruptedException {
return e.invokeAll(tasks, timeout, unit);
}
public <T> T invokeAny(Collection<? extends Callable<T>> tasks)
throws InterruptedException, ExecutionException {
return e.invokeAny(tasks);
}
public <T> T invokeAny(Collection<? extends Callable<T>> tasks,
long timeout, TimeUnit unit)
throws InterruptedException, ExecutionException, TimeoutException {
return e.invokeAny(tasks, timeout, unit);
}
}
ScheduledExecutorService
Предоставьте ряд методов расписания, чтобы задачи можно было откладывать или выполнять периодически. Соответствующий метод расписания будет возвращать ScheduledFuture, чтобы подтвердить, следует ли выполнять или отменять. Его реализация ScheduledThreadPoolExecutor также поддерживает немедленное выполнение задач, отправленных с помощью submit.
Поддерживается только относительное время задержки, например выполнение через 5 минут. Подобный таймер также может управлять отложенными задачами и периодическими задачами, но есть некоторые недостатки:
- Все задачи синхронизации имеют только один поток. Если задача выполняется долго, это повлияет на точность других задач TimerTask.
ScheduledExecutorService的多线程机制可弥补- TimerTask генерирует непроверенное исключение, которое прерывает выполнение потока и ошибочно считает задачу отмененной.
1:可以使用try-catch-finally对相应执行快处理;2:通过execute执行的方法可以设置UncaughtExceptionHandler来接收未捕获的异常,并作出处理;3:通过submit执行的,将被封装层ExecutionException重新抛出
ThreadPoolExecutor
public ThreadPoolExecutor(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue,
ThreadFactory threadFactory,
RejectedExecutionHandler handler) {
if (corePoolSize < 0 ||
maximumPoolSize <= 0 ||
maximumPoolSize < corePoolSize ||
keepAliveTime < 0)
throw new IllegalArgumentException();
if (workQueue == null || threadFactory == null || handler == null)
throw new NullPointerException();
this.corePoolSize = corePoolSize;
this.maximumPoolSize = maximumPoolSize;
this.workQueue = workQueue;
this.keepAliveTime = unit.toNanos(keepAliveTime);
this.threadFactory = threadFactory;
this.handler = handler;
}
- corePoolSize, maxPoolSize: ThreadPoolExecutor автоматически настраивает размер пула потоков в соответствии с этими двумя, когда новая задача отправляется через выполнение:
Если количество запущенных в данный момент потоков меньше, чем corePoolSize, создайте новый поток;
Если текущее количество потоков находится между corePoolSize и maxPoolSize, новые потоки будут создаваться только при заполнении очереди;
Если достигнуто максимальное количество потоков и очередь заполнена, политика отклонения будет выполняться в этом насыщенном состоянии.По умолчанию поток будет запускаться только при поступлении новой задачи, которую можно запустить заранее с помощью метода prestartCoreThread.
- corePoolSize: минимальное количество рабочих процессов, которое должен поддерживать пул потоков по умолчанию, даже если срок действия рабочих процессов истечет. Если вы не хотите его сохранять, вам нужно установить allowCoreThreadTimeOut, а наименьшее значение равно 0.
- maxPoolSize: максимальное количество потоков в пуле потоков. Предел Java составляет до 2 ^ 29-1, около 500 миллионов
- keepAliveTime, единица: если в текущем пуле потоков больше потоков, чем corePoolSize, пока время простоя любого потока превышает настройку keepAliveTime, он будет завершен; единица — это его единица времени.
- workQueue: можно использовать любую BlockingQueue, в основном их три.
- Прямая передача, прямая доставка задач. Например, SynchronousQueue, если нет потребления потока, отправка задачи не удастся, конечно, для ее обработки можно создать новый поток. Он подходит для обработки задач с зависимостями, и, как правило, его максимальный размер пула будет равен наибольшему.
- Неограниченные очереди, неограниченные очереди. Например, LinkedBlockingQueue, что означает, что если выполняются потоки corePoolSize, другие задачи могут только ждать. Он подходит для обработки задач, которые не зависят друг от друга,
- Ограниченные очереди, ограниченные очереди. Например, ArrayBlockingQueue необходимо учитывать взаимосвязь между размером очереди и максимальным количеством потоков, чтобы добиться лучшего использования ресурсов и пропускной способности.
- threadFactory: если не указано, используйте
Executors.defaultThreadFactory - RejectedExecutionHandler: задача, добавленная с помощью execute, будет выполнена, если Executor был закрыт или насыщен (количество потоков достигает максимального размера пула, а очередь заполнена) Java предоставляет 4 стратегии:
- AbortPolicy, выдает исключение RejectedExecutionException во время выполнения при отклонении;
- CallerRunsPolicy, если исполнитель не закрыт, выполняется тем потоком, который вызвал execute;
- DiscardPolicy, сразу выбрасывать новые задачи;
- DiscardOldestPolicy, если экзекьютор не закрыт, то отбросить задачу во главе очереди и повторить попытку;
ThreadPoolExecutor можно настроить beforeExecutor, afterExecutor можно использовать для добавления статистики журнала, времени, событий или функций сбора статистики, независимо от того, завершается ли run нормально или выдает исключение, afterExecutor будет выполнен. Если beforeExecutor выдает исключение RuntimeException, ни задача, ни afterExecutor не будут выполнены. terminated вызывается после завершения всех задач и закрытия всех рабочих потоков.Его также можно использовать для отправки уведомлений, записи журналов и т. д. в это время.
Как оценить размер пула потоков
- Вычислительно интенсивный, обычно
В системах с несколькими процессорами размер пула потоков устанавливается равным
может достичь оптимального использования;
количество процессоров
- Интенсивный ввод-вывод или другие блокирующие задачи, определите
количество процессоров,
загрузка процессора,
— это отношение времени ожидания ко времени вычислений, а оптимальный размер пула потоков в это время равен
описание сцены
Разделите бизнес веб-сайта на следующие части
- Прием заявок от клиентов и обработка заявок
- Текст и изображения, возвращаемые рендерингом страницы
- Получить рекламу для страницы
Прием и обработка запросов
теоретическая модель
Теоретически сервер может получать запросы и обрабатывать непрерывные запросы, реализуя согласованный интерфейс.
ServerSocket socket = new ServerSocket(80);
while(true){
Socket conn = socket.accept();
handleRequest(conn)
}
Недостаток: единовременно может обрабатываться только один запрос.При поступлении нового запроса он должен дождаться завершения обработки обрабатываемого запроса перед получением нового запроса.
Показать Создать многопоточность
Создайте новый поток для обслуживания каждого запроса
ServerSocket socket = new ServerSocket(80);
while(true){
final Socket conn = socket.accept();
Runnable task = new Runnable(){
public void run(){
handleRequest(conn);
}
}
new Thread(task).start();
}
недостаток:
- Создание и уничтожение потоков имеют определенные накладные расходы, задерживающие обработку запросов;
- Создается больше потоков, чем доступно процессоров, в результате чего потоки простаивают, что может оказать давление на сборку мусора.
- Большое количество выживших потоков, конкурирующих за ресурсы ЦП, приведет к большим потерям производительности.
- Существует ограничение на количество потоков, которые могут быть созданы в системе.
Использовать пул потоков
Используйте среду Executor, которая поставляется с java.
private static final Executor exec = Executors.newFixedThreadPool(100);
...
ServerSocket socket = new ServerSocket(80);
while(true){
final Socket conn = socket.accept();
Runnable task = new Runnable(){
public void run(){
handleRequest(conn);
}
}
exec.execute(task);
}
...
Стратегия пула потоков ограничивает количество одновременных задач, повторно использует существующие потоки и решает проблемы исчерпания ресурсов, чрезмерной конкуренции и частого создания каждого потока, созданного за счет реализации предполагаемых требований к потоку, а также включает в себя преимущества потоков. подача и выполнение задания.
Текст и изображения, возвращаемые рендерингом страницы
серийный рендеринг
renderText(source);
List<ImageData> imageData = new ArrayList<ImageData>();
for(ImageInfo info:scaForImageInfo(source)){
imageData.add(info.downloadImage());
}
for(ImageData data:imageData){
renderImage(data);
}
Недостатки: Большую часть времени загрузки изображений приходится на ожидание завершения операции ввода/вывода, в этот период процессор почти не работает, что заставляет пользователя долго ждать, прежде чем он увидит финальную страницу.
Распараллеливание
Процесс рендеринга можно разделить на две части: 1 — визуализировать текст, 1 — загрузить изображение.
private static final ExecutorService exec = Executors.newFixedThreadPool(100);
...
final List<ImageInfo> infos=scaForImageInfo(source);
Callable<List<ImageData>> task=new Callable<List<ImageData>>(){
public List<ImageData> call(){
List<ImageData> r = new ArrayList<ImageData>();
for(ImageInfo info:infos){
r.add(info.downloadImage());
}
return r;
}
};
Future<List<ImageData>> future = exec.submit(task);
renderText(source);
try{
List<ImageData> imageData = future.get();
for(ImageData data:imageData){
renderImage(data);
}
}catch(InterruptedException e){
Thread.currentThread().interrupt();
future.cancel(true);
}catche(ExecutionException e){
throw launderThrowable(e.getCause());
}
Используйте Callable, чтобы вернуть результат загруженного изображения, и используйте future, чтобы получить загруженное изображение, что сократит время ожидания, требуемое пользователем.
Недостатки: Время загрузки картинок явно медленнее, чем текста, такое распараллеливание скорее всего увеличит скорость всего на 1%
Параллельное улучшение производительности
Воспользуйтесь сервисом завершения.
private static final ExecutorService exec;
...
final List<ImageInfo> infos=scaForImageInfo(source);
CompletionService<ImageData> cService = new ExecutorCompletionService<ImageData>(exec);
for(final ImageInfo info:infos){
cService.submit(new Callable<ImageData>(){
public ImageData call(){
return info.downloadImage();
}
});
}
renderText(source);
try{
for(int i=0,n=info.size();t<n;t++){
Future<ImageData> f = cService.take();
ImageData imageData=f.get();
renderImage(imageData)
}
}catch(InterruptedException e){
Thread.currentThread().interrupt();
}catche(ExecutionException e){
throw launderThrowable(e.getCause());
}
Основная идея состоит в том, чтобы создать независимую задачу для каждой загрузки изображения и выполнить их в пуле потоков, тем самым преобразовав процесс последовательной загрузки в параллельный процесс.
Получить рекламу для страницы
Если отображение рекламы не получено в течение определенного периода времени, оно больше не может отображаться, и задание тайм-аута может быть отменено.
ExecutorService exe = Executors.newFixedThreadPool(3);
List<MyTask> myTasks = new ArrayList<>();
for (int i=0;i<3;i++){
myTasks.add(new MyTask(3-i));
}
try {
List<Future<String>> futures = exe.invokeAll(myTasks, 1, TimeUnit.SECONDS);
for (int i=0;i<futures.size();i++){
try {
String s = futures.get(i).get();
System.out.println("task execut "+myTasks.get(i).getSleepSeconds()+" s");
} catch (ExecutionException e) {
System.out.println("task sleep "+myTasks.get(i).getSleepSeconds()+" not execute ");
}catch (CancellationException e){
System.out.println("task sleep "+myTasks.get(i).getSleepSeconds()+" not execute ,because "+e);
}
}
} catch (InterruptedException e) {
e.printStackTrace();
}
exe.shutdown();
Метод invokeAll будет отменен для незавершенных задач, которые могут быть перехвачены с помощью CancellationException.Порядок последовательности, возвращаемый invokeAll, соответствует входящей задаче. Результат выглядит следующим образом:
task sleep 3 not execute ,because java.util.concurrent.CancellationException
task sleep 2 not execute ,because java.util.concurrent.CancellationException
task execut 1 s