Класс SubmissionPublisher<T>

Параметры типа:
T - тип публикуемого элемента
Все реализуемые интерфейсы:
AutoCloseable, Flow.Publisher<T>
public class SubmissionPublisher<T>
extends Object
implements Flow.Publisher<T>, AutoCloseable

A Flow.Publisher который асинхронно выдает отправленные (не пустые) элементы текущим подписчикам до его закрытия. Каждый текущий подписчик получает новые отправленные элементы в том же порядке, если не возникают ошибки или исключения. Использование SubmissionPublisher позволяет генераторам элементов действовать как совместимые reactive-streams Publishers, полагаясь на обработку ошибок и/или блокировку для управления потоком.

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

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

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

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

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

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

Этот класс также может служить удобной основой для подклассов, которые генерируют элементы и используют методы этого класса для их публикации. Например, вот класс, который периодически публикует элементы, сгенерированные из поставщика. (На практике вы можете добавить методы для независимого запуска и остановки генерации, для совместного использования Executors между издателями и так далее, или использовать 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(); }
 }
С:
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, заданный обработчик, вызываемый, когда любой Subscriber вызывает исключение в методе 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, если заданный Subscriber в настоящее время подписан.

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)

Добавляет данного 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 - исполнители для асинхронной доставки, поддерживающие создание как минимум одной независимой потоковой нити
maxBufferCapacity - максимальная ёмкость буфера каждого подписчика (принудительная ёмкость может быть округлённа до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым этой реализацией; метод getMaxBufferCapacity() возвращает фактическое значение)
handler - если не null, процедура вызова при возникновении исключения в методе onNext
Исключения:
NullPointerException - если executor равен null
IllegalArgumentException - если maxBufferCapacity не положительное

SubmissionPublisher

public SubmissionPublisher(Executor executor,
                           int maxBufferCapacity)

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

Параметры:
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.

Specified by:
subscribe в интерфейсе Flow.Publisher<T>
Parameters:
subscriber - подписчик
Throws:
NullPointerException - если подписчик null

submit

public int submit(T item)

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

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

Parameters:
item - публикуемый (не null) элемент
Returns:
оценка максимальной задержки среди подписчиков
Throws:
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) при попытке асинхронного уведомления подписчиков, или обработчик отбрасывания вызывает исключение при обработке отброшенного элемента, то это исключение перебрасывается.

Parameters:
item - публикуемый (не null) элемент
onDrop - если не null, обработчик, вызываемый при отбрасывании подписчику, с аргументами подписчик и элемент; если он возвращает true, попытка предложения повторяется (один раз)
Returns:
если отрицательный, (отрицательное) число отброшенных; в противном случае оценка максимальной задержки
Throws:
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) при попытке асинхронного уведомления подписчиков, или обработчик отбрасывания вызывает исключение при обработке отброшенного элемента, то это исключение перебрасывается.

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

close

public void close()

Если не закрыт, отправляет сигналы onComplete текущим подписчикам и запрещает последующие попытки публикации. По возвращении, этот метод НЕ гарантирует, что все подписчики завершили работу.

Specified by:
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 отменяется, в этом случае дальнейшие элементы не обрабатываются.

Параметры:
consumer - функция, применяемая к каждому элементу onNext
Возвращает:
CompletableFuture, который завершается успешно, когда издатель сигнализирует об onComplete, и с исключением при любой ошибке или отмене
Исключения:
NullPointerException - если потребитель равен null

© 1993, 2020, 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/11/docs/api/java.base/java/util/concurrent/SubmissionPublisher.html

Spec-Zone .ru
спецификации, руководства, описания, API