Spec-Zone.ru › OpenJDK 24

Класс 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 that asynchronously issues submitted (non-null) items to current subscribers until it is closed. Each current subscriber receives newly submitted items in the same order unless drops or exceptions are encountered. Using a SubmissionPublisher allows item generators to act as compliant reactive-streams Publishers relying on drop handling and/or blocking for flow control.

A SubmissionPublisher uses the Executor supplied in its constructor for delivery to subscribers. The best choice of Executor depends on expected usage. If the generator(s) of submitted items run in separate threads, and the number of subscribers can be estimated, consider using a Executors.newFixedThreadPool(int). Otherwise consider using the default, normally the ForkJoinPool.commonPool().

Buffering allows producers and consumers to transiently operate at different rates. Each subscriber uses an independent buffer. Buffers are created upon first use and expanded as needed up to the given maximum. (The enforced capacity may be rounded up to the nearest power of two and/or bounded by the largest value supported by this implementation.) Invocations of request do not directly result in buffer expansion, but risk saturation if unfilled requests exceed the maximum capacity. The default value of Flow.defaultBufferSize() may provide a useful starting point for choosing a capacity based on expected rates, resources, and usages.

A single SubmissionPublisher may be shared among multiple sources. Actions in a source thread prior to publishing an item or issuing a signal happen-before actions subsequent to the corresponding access by each subscriber. But reported estimates of lag and demand are designed for use in monitoring, not for synchronization control, and may reflect stale or inaccurate views of progress.

Publication methods support different policies about what to do when buffers are saturated. Method submit blocks until resources are available. This is simplest, but least responsive. The offer methods may drop items (either immediately or with bounded timeout), but provide an opportunity to interpose a handler and then retry.

If any Subscriber method throws an exception, its subscription is cancelled. If a handler is supplied as a constructor argument, it is invoked before cancellation upon an exception in method onNext, but exceptions in methods onSubscribe, onError and onComplete are not recorded or handled before cancellation. If the supplied Executor throws RejectedExecutionException (or any other RuntimeException or Error) when attempting to execute a task, or a drop handler throws an exception when processing a dropped item, then the exception is rethrown. In these cases, not all subscribers will have been issued the published item. It is usually good practice to closeExceptionally in these cases.

Method consume(Consumer) simplifies support for a common case in which the only action of a subscriber is to request and process all items using a supplied function.

This class may also serve as a convenient base for subclasses that generate items, and use the methods in this class to publish them. For example here is a class that periodically publishes the items generated from a supplier. (In practice you might add methods to independently start and stop generation, to share Executors among publishers, and so on, or use a SubmissionPublisher as a component rather than a superclass.)

 
 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

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, если данный подписчик (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, 2025, 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://download.java.net/java/early_access/jdk24/docs/api/java.base/java/util/concurrent/SubmissionPublisher.html

Spec-Zone.ru

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