Spec-Zone.ru › OpenJDK 17

Класс SubmissionPublisher<T>

java.lang.Object
java.util.concurrent.SubmissionPublisher<T>
Type Parameters:
T - тип публикуемого элемента
Все реализуемые интерфейсы:
AutoCloseable, Flow.Publisher<T>
public class SubmissionPublisher<T> extends Object implements Flow.Publisher<T>, AutoCloseable
Класс Flow.Publisher, асинхронно выдающий переданные (не нулевые) элементы текущим подписчикам до момента закрытия. Каждый текущий подписчик получает новые переданные элементы в том же порядке, если не возникают сбои или исключения. Использование SubmissionPublisher позволяет генераторам элементов действовать как совместимые reactive-streams издатели, опираясь на обработку падения и/или блокировку для управления потоком.

SubmissionPublisher использует Executor, предоставленный в его конструкторе, для доставки подписчикам. Лучший выбор Executor зависит от ожидаемого использования. Если генераторы переданных элементов работают в отдельных потоках, и число подписчиков может быть оценено, рассмотрите использование Executors.newFixedThreadPool(int). В противном случае рассмотрите использование по умолчанию, обычно ForkJoinPool.commonPool().

Буферизация позволяет производителям и потребителям временно работать с разными скоростями. Каждый подписчик использует независимый буфер. Буферы создаются при первом использовании и расширяются по мере необходимости до заданного максимума. (Принудительная емкость может быть округлено до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым данной реализацией.) Вызовы request напрямую не приводят к расширению буфера, но есть риск переполнения, если необработанные запросы превысят максимальную емкость. Значение по умолчанию Flow.defaultBufferSize() может служить полезной отправной точкой для выбора емкости, основываясь на ожидаемых скоростях, ресурсах и использовании.

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

Методы публикации поддерживают разные политики относительно действий при переполнении буферов. Метод submit блокируется до тех пор, пока ресурсы не станут доступны. Это самый простой, но наименее быстрый способ. Методы offer могут отбрасывать элементы (немедленно или с ограниченным временем ожидания), но дают возможность вставить обработчик, а затем повторить попытку.

Если какой-либо метод подписчика бросает исключение, его подписка отменяется. Если в качестве аргумента конструктора предоставляется обработчик, он вызывается перед отменением при возникновении исключения в методе onNext, но исключения в методах onSubscribe, onError и onComplete не регистрируются и не обрабатываются до отмены. Если предоставленный Executor бросает RejectedExecutionException (или любое другое RuntimeException или Error) при попытке выполнить задачу, или обработчик отбрасывания бросает исключение при обработке отброшенного элемента, то исключение повторно выбрасывается. В этих случаях не все подписчики получат опубликованный элемент. В таких случаях обычно рекомендуется использовать closeExceptionally.

Метод consume(Consumer) упрощает поддержку распространенного случая, когда единственное действие подписчика — запросить и обработать все элементы с использованием предоставленной функции.

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

 
 class PeriodicPublisher<T> extends SubmissionPublisher<T> {
   final ScheduledFuture<?> periodicTask;
   final ScheduledExecutorService scheduler;
   PeriodicPublisher(Executor executor, int maxBufferCapacity,
                     Supplier<? extends T> supplier,
                     long period, TimeUnit unit) {
     super(executor, maxBufferCapacity);
     scheduler = new ScheduledThreadPoolExecutor(1);
     periodicTask = scheduler.scheduleAtFixedRate(
       () -> submit(supplier.get()), 0, period, unit);
   }
   public void close() {
     periodicTask.cancel(false);
     scheduler.shutdown();
     super.close();
   }
 }

Вот пример реализации Flow.Processor. Для простоты иллюстрации он использует запросы в один шаг своему издателю. Более адаптивная версия могла бы отслеживать поток, используя оценку задержки, возвращаемую submit, вместе с другими служебными методами.

 
 class TransformProcessor<S,T> extends SubmissionPublisher<T>
   implements Flow.Processor<S,T> {
   final Function<? super S, ? extends T> function;
   Flow.Subscription subscription;
   TransformProcessor(Executor executor, int maxBufferCapacity,
                      Function<? super S, ? extends T> function) {
     super(executor, maxBufferCapacity);
     this.function = function;
   }
   public void onSubscribe(Flow.Subscription subscription) {
     (this.subscription = subscription).request(1);
   }
   public void onNext(S item) {
     subscription.request(1);
     submit(function.apply(item));
   }
   public void onError(Throwable ex) { closeExceptionally(ex); }
   public void onComplete() { close(); }
 }
Since:
9

Краткое описание конструкторов

Конструктор Описание
SubmissionPublisher()
Создает новый SubmissionPublisher, используя ForkJoinPool.commonPool() для асинхронной доставки подписчикам (если он не поддерживает уровень параллелизма не менее двух, в этом случае создается новый поток для выполнения каждой задачи), с максимальной емкостью буфера Flow.defaultBufferSize() и без обработчика исключений Subscriber в методе onNext.
SubmissionPublisher(Executor executor, int maxBufferCapacity)
Создает новый SubmissionPublisher, используя указанный Executor для асинхронной доставки подписчикам, с заданным максимальным размером буфера для каждого подписчика и без обработчика исключений Subscriber в методе onNext.
SubmissionPublisher(Executor executor, int maxBufferCapacity, BiConsumer<? super Flow.Subscriber<? super T>,? super Throwable> handler)
Создает новый SubmissionPublisher, используя указанный Executor для асинхронной доставки подписчикам, с заданным максимальным размером буфера для каждого подписчика и, если не null, указанным обработчиком, вызываемым, когда любой подписчик выбрасывает исключение в методе onNext.

Краткое описание методов

Модификатор и тип Метод Описание
void close()
Если еще не закрыт, отправляет сигналы onComplete текущим подписчикам и запрещает последующие попытки публикации.
void closeExceptionally(Throwable error)
Если еще не закрыт, отправляет сигналы onError текущим подписчикам с заданной ошибкой и запрещает последующие попытки публикации.
CompletableFuture<Void> consume(Consumer<? super T> consumer)
Обрабатывает все опубликованные элементы с помощью заданной функции Consumer.
int estimateMaximumLag()
Возвращает оценку максимального количества элементов, произведенных, но еще не потребленных среди всех текущих подписчиков.
long estimateMinimumDemand()
Возвращает оценку минимального количества элементов, запрошенных (через request), но еще не произведенных, среди всех текущих подписчиков.
Throwable getClosedException()
Возвращает исключение, связанное с closeExceptionally, или null, если не закрыт или закрыт нормально.
Executor getExecutor()
Возвращает Executor, используемый для асинхронной доставки.
int getMaxBufferCapacity()
Возвращает максимальную емкость буфера на подписчика.
int getNumberOfSubscribers()
Возвращает количество текущих подписчиков.
List<Flow.Subscriber<? super T>> getSubscribers()
Возвращает список текущих подписчиков для мониторинга и отслеживания, а не для вызова методов Flow.Subscriber на подписчиках.
boolean hasSubscribers()
Возвращает true, если у этого издателя есть подписчики.
boolean isClosed()
Возвращает true, если этот издатель не принимает новые публикации.
boolean isSubscribed(Flow.Subscriber<? super T> subscriber)
Возвращает true, если данный подписчик в настоящее время подписан.
int offer(T item, long timeout, TimeUnit unit, BiPredicate<Flow.Subscriber<? super T>,? super T> onDrop)
Опубликовывает указанный элемент, если возможно, каждому текущему подписчику, асинхронно вызывая его метод onNext, блокируя, пока ресурсы для любой подписки недоступны, вплоть до заданного таймаута или пока поток вызывающего не будет прерван, в этом случае вызывается заданный обработчик (если не null), и если он вернет true, повторить попытку один раз.
int offer(T item, BiPredicate<Flow.Subscriber<? super T>,? super T> onDrop)
Опубликовывает указанный элемент, если возможно, каждому текущему подписчику, асинхронно вызывая его метод onNext.
int submit(T item)
Опубликовывает указанный элемент каждому текущему подписчику, асинхронно вызывая его метод onNext, непрерывно блокируя, пока ресурсы для любого подписчика недоступны.
void subscribe(Flow.Subscriber<? super T> subscriber)
Добавляет данного подписчика, если он еще не подписан.

Методы, объявленные в классе java.lang.Object

clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait

Подробное описание конструкторов

SubmissionPublisher

public SubmissionPublisher(Executor executor, int maxBufferCapacity, BiConsumer<? super Flow.Subscriber<? super T>,? super Throwable> handler)
Создаёт новый SubmissionPublisher, использующий указанный Executor для асинхронной доставки подписчикам, с заданным максимальным размером буфера для каждого подписчика, и, если не null, указанным обработчиком, вызываемым, когда любой подписчик бросает исключение в методе onNext.
Параметры:
executor - Executor для использования при асинхронной доставке, поддерживающий создание как минимум одной независимой нити
maxBufferCapacity - максимальная ёмкость буфера каждого подписчика (принудительная ёмкость может быть округлено до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым данным реализацией; метод getMaxBufferCapacity() возвращает фактическое значение)
handler - если не null, процедура, вызываемая при возникновении исключения в методе onNext
Исключения:
NullPointerException - если executor равен null
IllegalArgumentException - если maxBufferCapacity не положительно

SubmissionPublisher

public SubmissionPublisher(Executor executor, int maxBufferCapacity)
Создаёт новый SubmissionPublisher, использующий указанный Executor для асинхронной доставки подписчикам, с заданным максимальным размером буфера для каждого подписчика, и без обработчика исключений подписчиков в методе onNext.
Параметры:
executor - Executor для использования при асинхронной доставке, поддерживающий создание как минимум одной независимой нити
maxBufferCapacity - максимальная ёмкость буфера каждого подписчика (принудительная ёмкость может быть округлено до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым данным реализацией; метод getMaxBufferCapacity() возвращает фактическое значение)
Исключения:
NullPointerException - если executor равен null
IllegalArgumentException - если maxBufferCapacity не положительно

SubmissionPublisher

public SubmissionPublisher()
Создаёт новый SubmissionPublisher, использующий ForkJoinPool.commonPool() для асинхронной доставки подписчикам (если он не поддерживает уровень параллелизма как минимум два, в этом случае создаётся новая нить для выполнения каждой задачи), с максимальной ёмкостью буфера Flow.defaultBufferSize() и без обработчика исключений подписчиков в методе onNext.

Подробное описание методов

subscribe

public void subscribe(Flow.Subscriber<? super T> subscriber)
Добавляет указанного подписчика, если он ещё не подписан. Если уже подписан, метод onError подписчика вызывается на существующей подписке с IllegalStateException. В противном случае, при успехе, метод onSubscribe подписчика вызывается асинхронно с новой Flow.Subscription. Если onSubscribe бросает исключение, подписка отменяется. В противном случае, если этот SubmissionPublisher был закрыт с исключением, вызывается метод onError подписчика с соответствующим исключением, или если закрыт без исключения, вызывается метод onComplete подписчика. Подписчики могут включить получение элементов, вызвав метод request новой подписки, и могут отписаться, вызвав метод cancel.
Определено в:
subscribe в интерфейсе Flow.Publisher<T>
Параметры:
subscriber - подписчик
Исключения:
NullPointerException - если подписчик равен null

submit

public int submit(T item)
Опубликовывает заданный элемент для каждого текущего подписчика, асинхронно вызывая его метод onNext, блокируя не прерывимо, пока ресурсы для любого подписчика недоступны. Данный метод возвращает оценку максимальной задержки (количество элементов, отправленных, но ещё не потреблённых) среди всех текущих подписчиков. Это значение не менее одного (учитывая этот отправленный элемент), если есть подписчики, иначе ноль.

Если Executor для этого издателя выбрасывает RejectedExecutionException (или любое другое RuntimeException или Error) при попытке асинхронного уведомления подписчиков, то это исключение повторно выбрасывается, в таком случае не все подписчики получат этот элемент.

Параметры:
item - публикуемый (не null) элемент
Возвращает:
оценка максимальной задержки среди подписчиков
Исключения:
IllegalStateException - если закрыто
NullPointerException - если элемент равен null
RejectedExecutionException - если брошено Executor

offer

public int offer(T item, BiPredicate<Flow.Subscriber<? super T>,? super T> onDrop)
Опубликовывает заданный элемент, если возможно, для каждого текущего подписчика, асинхронно вызывая его метод onNext. Элемент может быть отброшен одним или несколькими подписчиками, если ресурсы исчерпаны, в этом случае вызывается указанный обработчик (если не null), и если он возвращает true, повторяется ещё раз. Другие вызовы методов в этом классе другими потоками блокируются, пока вызывается обработчик. Если восстановление не гарантировано, обычно доступны только варианты: регистрация ошибки и/или отправка сигнала onError подписчику.

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

Если Executor для этого издателя выбрасывает RejectedExecutionException (или любое другое RuntimeException или Error) при попытке асинхронного уведомления подписчиков, или обработчик отбрасывания выбрасывает исключение при обработке отброшенного элемента, то это исключение повторно выбрасывается.

Параметры:
item - публикуемый (не null) элемент
onDrop - если не null, обработчик, вызываемый при отбрасывании подписчику, с аргументами подписчика и элемента; если он возвращает true, попытка отправки повторяется (один раз)
Возвращает:
если отрицательно, то (отрицательное) число отбросов; в противном случае оценка максимальной задержки
Исключения:
IllegalStateException - если закрыто
NullPointerException - если элемент равен null
RejectedExecutionException - если брошено Executor

offer

public int offer(T item, long timeout, TimeUnit unit, BiPredicate<Flow.Subscriber<? super T>,? super T> onDrop)
Опубликовывает заданный элемент, если возможно, для каждого текущего подписчика, асинхронно вызывая его метод onNext, блокируя, пока ресурсы для любой подписки недоступны, вплоть до указанного таймаута или пока вызывающий поток не прерван, в этот момент вызывается указанный обработчик (если не null), и если он возвращает true, повторяется ещё раз. (Обработчик отбрасывания может различать таймауты и прерывания, проверяя, прерван ли текущий поток). Другие вызовы методов в этом классе другими потоками блокируются, пока вызывается обработчик. Если восстановление не гарантировано, обычно доступны только варианты: регистрация ошибки и/или отправка сигнала onError подписчику.

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

Если Executor для этого издателя выбрасывает RejectedExecutionException (или любое другое RuntimeException или Error) при попытке асинхронного уведомления подписчиков, или обработчик отбрасывания выбрасывает исключение при обработке отброшенного элемента, то это исключение повторно выбрасывается.

Параметры:
item - публикуемый (не null) элемент
timeout - сколько ждать ресурсов для любого подписчика, прежде чем отказаться, в единицах unit
unit - TimeUnit определяя, как интерпретировать параметр timeout
onDrop - если не null, обработчик, вызываемый при отбрасывании подписчику, с аргументами подписчика и элемента; если он возвращает true, попытка отправки повторяется (один раз)
Возвращает:
если отрицательно, то (отрицательное) число отбросов; в противном случае оценка максимальной задержки
Исключения:
IllegalStateException - если закрыто
NullPointerException - если элемент равен null
RejectedExecutionException - если брошено Executor

close

public void close()
Если ещё не закрыто, отправляет сигналы onComplete текущим подписчикам и запрещает дальнейшие попытки публикации. По возвращении, этот метод НЕ гарантирует, что все подписчики уже завершились.
Определено в:
close в интерфейсе AutoCloseable

closeExceptionally

public void closeExceptionally(Throwable error)
Если издатель ещё не закрыт, он отправляет сигнал onError текущим подписчикам с указанной ошибкой и запрещает дальнейшие попытки публикации. Будущие подписчики также получат указанную ошибку. По завершении метод НЕ гарантирует, что все подписчики завершили работу.
Parameters:
error - аргумент onError, отправленный подписчикам
Throws:
NullPointerException - если ошибка равна null

isClosed

public boolean isClosed()
Возвращает true, если этот издатель больше не принимает новые подписки.
Returns:
true, если закрыт

getClosedException

public Throwable getClosedException()
Возвращает исключение, связанное с closeExceptionally, или null, если издатель не закрыт или закрыт без ошибок.
Returns:
исключение или null, если отсутствует

hasSubscribers

public boolean hasSubscribers()
Возвращает true, если у этого издателя есть подписчики.
Returns:
true, если у издателя есть подписчики

getNumberOfSubscribers

public int getNumberOfSubscribers()
Возвращает количество текущих подписчиков.
Returns:
количество текущих подписчиков

getExecutor

public Executor getExecutor()
Возвращает Executor, используемый для асинхронной доставки.
Returns:
Executor, используемый для асинхронной доставки

getMaxBufferCapacity

public int getMaxBufferCapacity()
Возвращает максимальную емкость буфера на подписчика.
Returns:
максимальную емкость буфера на подписчика

getSubscribers

public List<Flow.Subscriber<? super T>> getSubscribers()
Возвращает список текущих подписчиков для мониторинга и отслеживания, а не для вызова методов Flow.Subscriber на подписчиках.
Returns:
список текущих подписчиков

isSubscribed

public boolean isSubscribed(Flow.Subscriber<? super T> subscriber)
Возвращает true, если указанный подписчик сейчас подписан.
Parameters:
subscriber - подписчик
Returns:
true, если подписан
Throws:
NullPointerException - если подписчик равен null

estimateMinimumDemand

public long estimateMinimumDemand()
Возвращает оценку минимального количества запрошенных элементов (через request), но еще не произведенных, среди всех текущих подписчиков.
Returns:
оценку или ноль, если нет подписчиков

estimateMaximumLag

public int estimateMaximumLag()
Возвращает оценку максимального количества произведенных, но еще не потребленных элементов среди всех текущих подписчиков.
Returns:
оценку

consume

public CompletableFuture<Void> consume(Consumer<? super T> consumer)
Обрабатывает все опубликованные элементы с помощью заданной функции Consumer. Возвращает CompletableFuture, которая завершается нормально, когда этот издатель сигнализирует onComplete, или завершается с исключением при любой ошибке, или исключение выбрасывается функцией Consumer, или возвращаемая CompletableFuture отменяется, в этом случае дальнейшая обработка элементов не происходит.
Parameters:
consumer - функция, применяемая к каждому элементу onNext
Returns:
CompletableFuture, которая завершается нормально, когда издатель сигнализирует onComplete, и с исключением при любой ошибке или отмене
Throws:
NullPointerException - если consumer равен null

© 1993, 2021, Oracle and/or its affiliates. All rights reserved.
Documentation extracted from Debian's OpenJDK Development Kit package.
Licensed under the GNU General Public License, version 2, with the Classpath Exception.
Various third party code in OpenJDK is licensed under different licenses (see Debian package).
Java and OpenJDK are trademarks or registered trademarks of Oracle and/or its affiliates.
https://docs.oracle.com/en/java/javase/17/docs/api/java.base/java/util/concurrent/SubmissionPublisher.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API