Класс 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 может быть совместно использован несколькими источниками. Действия в потоке источника до публикации элемента или выдачи сигнала происходят раньше действий, последующих за соответствующим доступом каждого подписчика. Но заявленные оценки запаздывания и спроса предназначены для мониторинга, а не для управления синхронизацией, и могут отражать устаревшие или неточные взгляды на прогресс.
Методы публикации поддерживают различные стратегии действий при переполнении буферов. Метод 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, заданный обработчик, вызываемый при возникновении исключения у любого 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
closeИсключительно
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, 2023, 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/21/docs/api/java.base/java/util/concurrent/SubmissionPublisher.html