Spec-Zone.ru › OpenJDK 21

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

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

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

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

Методы публикации поддерживают различные стратегии действий при переполнении буферов. Метод 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, заданный обработчик, вызываемый при возникновении исключения у любого 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.
Определено в:
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

closeИсключительно

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, если указанный Subscriber сейчас подписан.
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, 2023, 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/21/docs/api/java.base/java/util/concurrent/SubmissionPublisher.html

Spec-Zone.ru

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