Класс SubmissionPublisher<T>
- Type Parameters:
T- тип публикуемого элемента
- Все реализованные интерфейсы:
-
AutoCloseable,Flow.Publisher<T>
public class SubmissionPublisher<T> extends Object implements Flow.Publisher<T>, AutoCloseable
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 |
Создаёт новый SubmissionPublisher, используя указанный Executor для асинхронной доставки подписчикам, с заданным максимальным размером буфера для каждого подписчика и без обработчика исключений Subscriber в методе onNext. |
SubmissionPublisher |
Создаёт новый SubmissionPublisher, используя указанный Executor для асинхронной доставки подписчикам, с заданным максимальным размером буфера для каждого подписчика и, если не null, заданным обработчиком, вызываемым, когда любой Subscriber выбрасывает исключение в методе onNext. |
Краткое описание методов
| Модификатор и тип | Метод | Описание |
|---|---|---|
void |
close() |
Если ещё не закрыт, издаёт сигналы onComplete текущим подписчикам и запрещает последующие попытки публикации. |
void |
closeExceptionally |
Если ещё не закрыт, издаёт сигналы onError текущим подписчикам с указанной ошибкой и запрещает последующие попытки публикации. |
CompletableFuture |
consume |
Обрабатывает все опубликованные элементы, используя заданную функцию Consumer. |
int |
estimateMaximumLag() |
Возвращает оценку максимального количества элементов, произведённых, но ещё не потреблённых среди всех текущих подписчиков. |
long |
estimateMinimumDemand() |
Возвращает оценку минимального количества элементов, запрошенных (через request), но ещё не произведённых, среди всех текущих подписчиков. |
Throwable |
getClosedException() |
Возвращает исключение, связанное с closeExceptionally, или null, если не закрыт или закрыт нормально. |
Executor |
getExecutor() |
Возвращает Executor, используемый для асинхронной доставки. |
int |
getMaxBufferCapacity() |
Возвращает максимальную ёмкость буфера на подписчика. |
int |
getNumberOfSubscribers() |
Возвращает количество текущих подписчиков. |
List |
getSubscribers() |
Возвращает список текущих подписчиков для целей мониторинга и отслеживания, а не для вызова методов Flow.Subscriber на подписчиках. |
boolean |
hasSubscribers() |
Возвращает true, если у этого издателя есть подписчики. |
boolean |
isClosed() |
Возвращает true, если этот издатель не принимает новые задания. |
boolean |
isSubscribed |
Возвращает true, если данный Subscriber в настоящее время подписан. |
int |
offer |
Опубликовывает данный элемент, если возможно, каждому текущему подписчику, асинхронно вызывая его метод onNext, блокируя, пока ресурсы для любой подписки недоступны, до указанного таймаута или пока поток вызывающей стороны не будет прерван, в этот момент вызывается заданный обработчик (если он не null), и если он возвращает true, повторяет попытку один раз. |
int |
offer |
Опубликовывает данный элемент, если возможно, каждому текущему подписчику, асинхронно вызывая его метод onNext. |
int |
submit |
Опубликовывает данный элемент каждому текущему подписчику, асинхронно вызывая его метод onNext, непрерывно блокируя, пока ресурсы для любого подписчика недоступны. |
void |
subscribe |
Добавляет данный Subscriber, если он ещё не подписан. |
Подробное описание конструкторов
SubmissionPublisher
public SubmissionPublisher(Executor executor, int maxBufferCapacity, BiConsumer<? super Flow.Subscriber<? super T>, ? super Throwable> handler)
onNext.- Параметры:
-
executor- исполнитель, который используется для асинхронной доставки, поддерживающий создание как минимум одной независимой потоковой нити -
maxBufferCapacity- максимальная ёмкость буфера каждого подписчика (применяемая ёмкость может быть округлено до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым данной реализацией; методgetMaxBufferCapacity()возвращает фактическое значение) -
handler- если не null, процедура, вызываемая при возникновении исключения в методеonNext - Исключения:
-
NullPointerException- если executor равен null -
IllegalArgumentException- если maxBufferCapacity не положительно
SubmissionPublisher
public SubmissionPublisher(Executor executor, int maxBufferCapacity)
onNext.- Параметры:
-
executor- исполнитель, который используется для асинхронной доставки, поддерживающий создание как минимум одной независимой потоковой нити -
maxBufferCapacity- максимальная ёмкость буфера каждого подписчика (применяемая ёмкость может быть округлено до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым данной реализацией; методgetMaxBufferCapacity()возвращает фактическое значение) - Исключения:
-
NullPointerException- если executor равен null -
IllegalArgumentException- если maxBufferCapacity не положительно
SubmissionPublisher
public 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()
- Returns:
- true, если закрыт
getClosedException
public Throwable getClosedException()
closeExceptionally, или null, если издатель не закрыт или закрыт нормально.- Returns:
- исключение, или null, если нет
hasSubscribers
public boolean hasSubscribers()
- Returns:
- true, если у издателя есть подписчики
getNumberOfSubscribers
public int getNumberOfSubscribers()
- Returns:
- количество текущих подписчиков
getExecutor
public Executor getExecutor()
- 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)
- 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)
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