Класс SubmissionPublisher<T>
- Type Parameters:
-
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 может быть общим для нескольких источников. Действия в потоке источника перед публикацией элемента или вызовом сигнала happen-before действия, последующие за соответствующим доступом каждого подписчика. Однако указанные оценки задержки и спроса предназначены для мониторинга, а не для управления синхронизацией и могут отражать устаревшие или неточные представления о прогрессе.
Методы публикации поддерживают разные политики относительно действий при переполнении буферов. Метод 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 |
Создает новый SubmissionPublisher, используя указанный Executor для асинхронной доставки подписчикам, с заданным максимальным размером буфера для каждого подписчика и без обработчика исключений Subscriber в методе onNext. |
SubmissionPublisher |
Создает новый SubmissionPublisher, используя указанный Executor для асинхронной доставки подписчикам, с заданным максимальным размером буфера для каждого подписчика и, если не null, указанным обработчиком, вызываемым, когда любой подписчик выбрасывает исключение в методе onNext. |
Краткое описание методов
| Модификатор и тип | Метод | Описание |
|---|---|---|
void |
close() |
Если еще не закрыт, отправляет сигналы onComplete текущим подписчикам и запрещает последующие попытки публикации. |
void |
closeExceptionally |
Если еще не закрыт, отправляет сигналы onError текущим подписчикам с заданной ошибкой и запрещает последующие попытки публикации. |
CompletableFuture<Void> |
consume |
Обрабатывает все опубликованные элементы с помощью заданной функции 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 |
Возвращает true, если данный подписчик в настоящее время подписан. |
int |
offer |
Опубликовывает указанный элемент, если возможно, каждому текущему подписчику, асинхронно вызывая его метод onNext, блокируя, пока ресурсы для любой подписки недоступны, вплоть до заданного таймаута или пока поток вызывающего не будет прерван, в этом случае вызывается заданный обработчик (если не null), и если он вернет true, повторить попытку один раз. |
int |
offer |
Опубликовывает указанный элемент, если возможно, каждому текущему подписчику, асинхронно вызывая его метод onNext. |
int |
submit |
Опубликовывает указанный элемент каждому текущему подписчику, асинхронно вызывая его метод onNext, непрерывно блокируя, пока ресурсы для любого подписчика недоступны. |
void |
subscribe |
Добавляет данного подписчика, если он еще не подписан. |
Подробное описание конструкторов
SubmissionPublisher
public SubmissionPublisher(Executor executor, int maxBufferCapacity, BiConsumer<? super Flow.Subscriber<? super T>,? super Throwable> handler)
onNext.- Параметры:
-
executor- Executor для использования при асинхронной доставке, поддерживающий создание как минимум одной независимой нити -
maxBufferCapacity- максимальная ёмкость буфера каждого подписчика (принудительная ёмкость может быть округлено до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым данным реализацией; методgetMaxBufferCapacity()возвращает фактическое значение) -
handler- если не null, процедура, вызываемая при возникновении исключения в методеonNext - Исключения:
-
NullPointerException- если executor равен null -
IllegalArgumentException- если maxBufferCapacity не положительно
SubmissionPublisher
public SubmissionPublisher(Executor executor, int maxBufferCapacity)
onNext.- Параметры:
-
executor- 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, 2021, 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/17/docs/api/java.base/java/util/concurrent/SubmissionPublisher.html