Класс SubmissionPublisher<T>
- Параметры типа:
T- тип публикуемого элемента
- Все реализуемые интерфейсы:
AutoCloseable, Flow.Publisher<T>
public class SubmissionPublisher<T> extends Object implements Flow.Publisher<T>, AutoCloseable
Flow.Publisher, который асинхронно передает отправленные (ненулевые) элементы текущим подписчикам до закрытия. Каждый текущий подписчик получает вновь отправленные элементы в том же порядке, если не происходит отбрасывание элементов или исключений. Использование SubmissionPublisher позволяет генераторам элементов выступать в роли совместимых издателей reactive-streams, полагающихся на обработку отбрасывания и/или блокировку для управления потоком. 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) упрощает поддержку распространенного случая, когда единственное действие подписчика — запрашивать и обрабатывать все элементы с помощью переданной функции.
Этот класс также может служить удобным базовым классом для подклассов, генерирующих элементы и использующих методы этого класса для их публикации. Например, ниже приведен класс, который периодически публикует элементы, созданные поставщиком. (На практике можно добавить методы для независимого запуска и остановки генерации, совместного использования 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 |
Создает новый 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() и без обработчика исключений Subscriber в методе onNext.Подробное описание методов
subscribe
public void subscribe(Flow.Subscriber<? super T> subscriber)
onError этого Subscriber вызывается для существующей подписки с исключением IllegalStateException. В противном случае при успешном выполнении метод onSubscribe Subscriber асинхронно вызывается с новой 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- публикуемый (ненулевой) элемент - Возвращает:
- оценку максимальной задержки среди подписчиков
- Вызывает:
-
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- публикуемый (ненулевой) элемент -
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- публикуемый (ненулевой) элемент -
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.
https://docs.oracle.com/en/java/javase/25/docs/api/java.base/java/util/concurrent/SubmissionPublisher.html