Класс SubmissionPublisher<T>
- Параметры типа:
T— тип публикуемого элемента
- Все реализуемые интерфейсы:
AutoCloseable, Flow.Publisher<T>
public class SubmissionPublisher<T> extends Object implements Flow.Publisher<T>, AutoCloseable
Flow.Publisher, асинхронно отправляющий текущим подписчикам переданные (не null) элементы до своего закрытия. Каждый текущий подписчик получает вновь переданные элементы в том же порядке, если не возникает потери элементов или исключений. Использование SubmissionPublisher позволяет генераторам элементов выступать в роли совместимых издателей реактивных потоков, использующих обработку потерь и/или блокировку для управления потоком. SubmissionPublisher использует Executor, переданный конструктору, для доставки элементов подписчикам. Выбор Executor зависит от предполагаемого использования. Если генераторы передаваемых элементов работают в отдельных потоках, а число подписчиков можно оценить, рассмотрите возможность использования Executors.newFixedThreadPool(int). В противном случае рассмотрите стандартный вариант — обычно это ForkJoinPool.commonPool().
Буферизация позволяет производителям и потребителям временно работать с разной скоростью. Каждый подписчик использует независимый буфер. Буферы создаются при первом использовании и при необходимости увеличиваются до заданного максимума. (Фактическая ёмкость может быть округлена вверх до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым этой реализацией.) Вызовы request напрямую не приводят к увеличению буфера, но могут вызвать его переполнение, если число неудовлетворённых запросов превысит максимальную ёмкость. Значение по умолчанию Flow.defaultBufferSize() может стать полезной отправной точкой при выборе ёмкости с учётом ожидаемых скоростей, ресурсов и сценариев использования.
Один экземпляр SubmissionPublisher можно использовать совместно из нескольких источников. Действия в потоке источника, предшествующие публикации элемента или отправке сигнала, происходят до действий, следующих за соответствующим обращением каждого подписчика. Однако сообщаемые оценки отставания и спроса предназначены для мониторинга, а не для управления синхронизацией, и могут отражать устаревшие или неточные данные о ходе выполнения.
Методы публикации поддерживают различные политики обработки переполнения буферов. Метод submit блокируется до появления доступных ресурсов. Это самый простой, но наименее отзывчивый вариант. Методы offer могут отбрасывать элементы (сразу или по истечении ограниченного времени ожидания), но позволяют вызвать обработчик и затем повторить попытку.
Если метод Subscriber выбрасывает исключение, его подписка отменяется. Если обработчик передан аргументом конструктора, он вызывается перед отменой подписки при исключении в методе 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) {
submit(function.apply(item));
subscription.request(1);
}
public void onError(Throwable ex) { closeExceptionally(ex); }
public void onComplete() { close(); }
}
- Начиная с версии:
- 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, если он ещё не подписан. |
Методы, объявленные в классе Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait | Модификатор и тип | Метод | Описание |
|---|---|---|
protected Object |
clone() |
Создаёт и возвращает копию этого объекта. |
boolean |
equals |
Указывает, равен ли этот объект какому-либо другому объекту. |
protected void |
finalize() |
Устарело, будет удалено: этот элемент API подлежит удалению в будущей версии. Финализация объявлена устаревшей и подлежит удалению в одном из будущих выпусков. |
final Class |
getClass() |
Возвращает класс времени выполнения этого Object. |
int |
hashCode() |
Возвращает хеш-код этого объекта. |
final void |
notify() |
Пробуждает один поток, ожидающий на мониторе этого объекта. |
final void |
notifyAll() |
Пробуждает все потоки, ожидающие на мониторе этого объекта. |
String |
toString() |
Возвращает строковое представление объекта. |
final void |
wait() |
Заставляет текущий поток ожидать пробуждения, обычно вследствие вызова notify или interrupt. |
final void |
wait |
Заставляет текущий поток ожидать пробуждения, обычно вследствие вызова notify или interrupt, либо до истечения заданного промежутка реального времени. |
final void |
wait |
Заставляет текущий поток ожидать пробуждения, обычно вследствие вызова notify или interrupt, либо до истечения заданного промежутка реального времени. |
Подробное описание конструкторов
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() и без обработчика исключений Subscriber в методе onNext.Подробное описание методов
subscribe
public void subscribe(Flow.Subscriber<? super T> subscriber)
onError вызывается для существующей подписки с исключением IllegalStateException. В противном случае при успешном выполнении метод onSubscribe подписчика асинхронно вызывается с новой Flow.Subscription. Если onSubscribe выбрасывает исключение, подписка отменяется. В противном случае, если SubmissionPublisher был закрыт с исключением, вызывается метод onError подписчика с соответствующим исключением; если же он был закрыт без исключения, вызывается метод onComplete. Подписчики могут разрешить получение элементов, вызвав метод request новой Subscription, и отказаться от подписки, вызвав её метод cancel.- Определено в:
-
subscribeв интерфейсеFlow.Publisher<T> - Параметры:
-
subscriber— подписчик - Исключения:
-
NullPointerException— если subscriber равен null
submit
public int submit(T item)
onNext, и без возможности прерывания блокируется, пока ресурсы для любого подписчика недоступны. Этот метод возвращает оценку максимального отставания (числа переданных, но ещё не потреблённых элементов) среди всех текущих подписчиков. Если есть подписчики, значение не меньше единицы (с учётом передаваемого элемента), иначе оно равно нулю. Если Executor этого издателя выбрасывает RejectedExecutionException (или любое другое RuntimeException либо Error) при попытке асинхронно уведомить подписчиков, исключение повторно выбрасывается; в этом случае элемент может быть отправлен не всем подписчикам.
- Параметры:
-
item— публикуемый элемент (не равен null) - Возвращает:
- оценку максимального отставания среди подписчиков
- Исключения:
-
IllegalStateException— если издатель закрыт -
NullPointerException— если item равен 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, выполняется одна повторная попытка offer - Возвращает:
- если значение отрицательное — отрицательное число потерь; в противном случае оценку максимального отставания
- Исключения:
-
IllegalStateException— если издатель закрыт -
NullPointerException— если item равен 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, выполняется одна повторная попытка offer - Возвращает:
- если значение отрицательное — отрицательное число потерь; в противном случае оценку максимального отставания
- Исключения:
-
IllegalStateException— если издатель закрыт -
NullPointerException— если item равен null -
RejectedExecutionException— если выброшено Executor
close
public void close()
onComplete и запрещает последующие попытки публикации. Для обеспечения одинакового порядка для всех подписчиков этот метод может ожидать завершения выполняющихся предложений. По возвращении метод НЕ гарантирует, что все подписчики уже завершили обработку.- Определено в:
-
closeв интерфейсеAutoCloseable
closeExceptionally
public void closeExceptionally(Throwable error)
onError с указанной ошибкой и запрещает последующие попытки публикации. Будущие подписчики также получат указанную ошибку. По возвращении метод НЕ гарантирует, что все подписчики уже завершили обработку.- Параметры:
-
error— аргументonError, отправляемый подписчикам - Исключения:
-
NullPointerException— если error равен null
isClosed
public boolean isClosed()
- Возвращает:
- true, если издатель закрыт
getClosedException
public Throwable getClosedException()
closeExceptionally, или null, если издатель не закрыт либо закрыт штатно.- Возвращает:
- исключение или null, если его нет
hasSubscribers
public boolean hasSubscribers()
- Возвращает:
- true, если у этого издателя есть подписчики
getNumberOfSubscribers
public int getNumberOfSubscribers()
- Возвращает:
- число текущих подписчиков
getExecutor
public Executor getExecutor()
- Возвращает:
- Executor, используемый для асинхронной доставки
getMaxBufferCapacity
public int getMaxBufferCapacity()
- Возвращает:
- максимальную ёмкость буфера для каждого подписчика
getSubscribers
public List<Flow.Subscriber<? super T>> getSubscribers()
Flow.Subscriber у подписчиков.- Возвращает:
- список текущих подписчиков
isSubscribed
public boolean isSubscribed(Flow.Subscriber<? super T> subscriber)
- Параметры:
-
subscriber— подписчик - Возвращает:
- true, если подписка активна
- Исключения:
-
NullPointerException— если subscriber равен null
estimateMinimumDemand
public long estimateMinimumDemand()
request), но ещё не произведённых всеми текущими подписчиками.- Возвращает:
- оценку или ноль, если подписчиков нет
estimateMaximumLag
public int estimateMaximumLag()
- Возвращает:
- оценку
consume
public CompletableFuture<Void> consume(Consumer<? super T> consumer)
onComplete, или завершается исключением при любой ошибке, при выбрасывании исключения Consumer либо при отмене возвращённого CompletableFuture; в последнем случае дальнейшая обработка элементов не выполняется.- Параметры:
-
consumer— функция, применяемая к каждому элементу onNext - Возвращает:
- CompletableFuture, который завершается штатно при отправке издателем сигнала onComplete и завершается исключением при любой ошибке или отмене
- Исключения:
-
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.