Класс SubmissionPublisher<T>
- java.lang.Object
-
- java.util.concurrent.SubmissionPublisher<T>
- Параметры типа:
-
T- тип публикуемого элемента
- Все реализуемые интерфейсы:
-
AutoCloseable,Flow.Publisher<T>
public class SubmissionPublisher<T> extends Object implements Flow.Publisher<T>, AutoCloseable
A Flow.Publisher который асинхронно выдает отправленные (не пустые) элементы текущим подписчикам до его закрытия. Каждый текущий подписчик получает новые отправленные элементы в том же порядке, если не возникают ошибки или исключения. Использование SubmissionPublisher позволяет генераторам элементов действовать как совместимые reactive-streams Publishers, полагаясь на обработку ошибок и/или блокировку для управления потоком.
SubmissionPublisher использует Executor, предоставленный в его конструкторе для доставки подписчикам. Лучший выбор Executor зависит от ожидаемого использования. Если генератор(ы) отправляемых элементов работают в отдельных потоках, и количество подписчиков можно оценить, рассмотрите использование Executors.newFixedThreadPool(int). В противном случае рассмотрите использование по умолчанию, обычно ForkJoinPool.commonPool().
Буферизация позволяет производителям и потребителям временно работать с разными скоростями. Каждый подписчик использует независимый буфер. Буферы создаются при первом использовании и расширяются по мере необходимости до заданного максимума. (Принудительная емкость может быть округлена до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым этой реализацией.) Вызовы request не напрямую приводят к расширению буфера, но рискуют насыщением, если необработанные запросы превысят максимальную емкость. Значение по умолчанию Flow.defaultBufferSize() может служить хорошей отправной точкой для выбора емкости на основе ожидаемых скоростей, ресурсов и использования.
Один SubmissionPublisher может быть разделен между несколькими источниками. Действия в потоке источника перед публикацией элемента или выдачей сигнала happen-before последующие действия соответствующего доступа каждого подписчика. Однако указанные оценки задержки и спроса предназначены для мониторинга, а не для управления синхронизацией и могут отражать устаревшие или неточные представления о ходе работы.
Методы публикации поддерживают различные политики относительно того, что делать, когда буферы заполнены. Метод 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, использующий |
SubmissionPublisher(Executor executor,
int maxBufferCapacity) | Создает новый SubmissionPublisher, использующий предоставленный Executor для асинхронной доставки подписчикам, с заданным максимальным размером буфера для каждого подписчика и без обработчика исключений Subscriber в методе |
SubmissionPublisher(Executor executor,
int maxBufferCapacity,
BiConsumer<? super Flow.Subscriber<? super T>,? super Throwable> handler) | Создает новый SubmissionPublisher, использующий предоставленный Executor для асинхронной доставки подписчикам, с заданным максимальным размером буфера для каждого подписчика и, если не null, заданный обработчик, вызываемый, когда любой Subscriber вызывает исключение в методе |
Методы
| Модификатор и тип | Метод | Описание |
|---|---|---|
void | close() | Если не закрыт, отправляет сигналы |
void | closeExceptionally(Throwable error) | Если не закрыт, отправляет сигналы |
CompletableFuture<Void> | consume(Consumer<? super T> consumer) | Обрабатывает все опубликованные элементы с использованием заданной функции Consumer. |
int | estimateMaximumLag() | Возвращает оценку максимального количества произведённых, но ещё не потреблённых элементов среди всех текущих подписчиков. |
long | estimateMinimumDemand() | Возвращает оценку минимального количества запрошенных (через |
Throwable | getClosedException() | Возвращает исключение, связанное с |
Executor | getExecutor() | Возвращает Executor, используемый для асинхронной доставки. |
int | getMaxBufferCapacity() | Возвращает максимальную ёмкость буфера на подписчика. |
int | getNumberOfSubscribers() | Возвращает количество текущих подписчиков. |
List<Flow.Subscriber<? super T>> | getSubscribers() | Возвращает список текущих подписчиков для целей мониторинга и отслеживания, а не для вызова методов |
boolean | hasSubscribers() | Возвращает true, если у этого издателя есть подписчики. |
boolean | isClosed() | Возвращает true, если этот издатель не принимает отправки. |
boolean | isSubscribed(Flow.Subscriber<? super T> subscriber) | Возвращает true, если заданный Subscriber в настоящее время подписан. |
int | offer(T item,
long timeout,
TimeUnit unit,
BiPredicate<Flow.Subscriber<? super T>,? super T> onDrop) | Публикует данный элемент, если возможно, для каждого текущего подписчика, асинхронно вызывая его метод |
int | offer(T item,
BiPredicate<Flow.Subscriber<? super T>,? super T> onDrop) | Публикует данный элемент, если возможно, для каждого текущего подписчика, асинхронно вызывая его метод |
int | submit(T item) | Публикует данный элемент для каждого текущего подписчика, асинхронно вызывая его метод |
void | subscribe(Flow.Subscriber<? super T> subscriber) | Добавляет данного Subscriber, если он ещё не подписан. |
Методы, объявленные в классе java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait Конструкторы
SubmissionPublisher
public SubmissionPublisher(Executor executor,
int maxBufferCapacity,
BiConsumer<? super Flow.Subscriber<? super T>,? super Throwable> handler) Создаёт новый SubmissionPublisher, используя указанный Executor для асинхронной доставки подписчикам, с заданным максимальным размером буфера для каждого подписчика и, если не null, заданным обработчиком, вызываемым, когда любой подписчик выбрасывает исключение в методе onNext.
- Параметры:
-
executor- исполнители для асинхронной доставки, поддерживающие создание как минимум одной независимой потоковой нити -
maxBufferCapacity- максимальная ёмкость буфера каждого подписчика (принудительная ёмкость может быть округлённа до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым этой реализацией; методgetMaxBufferCapacity()возвращает фактическое значение) -
handler- если не null, процедура вызова при возникновении исключения в методеonNext - Исключения:
-
NullPointerException- если executor равен null -
IllegalArgumentException- если maxBufferCapacity не положительное
SubmissionPublisher
public SubmissionPublisher(Executor executor,
int maxBufferCapacity) Создаёт новый SubmissionPublisher, используя указанный Executor для асинхронной доставки подписчикам, с заданным максимальным размером буфера для каждого подписчика и без обработчика исключений подписчиков в методе onNext.
- Параметры:
-
executor- исполнители для асинхронной доставки, поддерживающие создание как минимум одной независимой потоковой нити -
maxBufferCapacity- максимальная ёмкость буфера каждого подписчика (принудительная ёмкость может быть округлённа до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым этой реализацией; методgetMaxBufferCapacity()возвращает фактическое значение) - Исключения:
-
NullPointerException- если executor равен null -
IllegalArgumentException- если maxBufferCapacity не положительное
SubmissionPublisher
public SubmissionPublisher()
Создаёт новый 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.
- Specified by:
-
subscribeв интерфейсеFlow.Publisher<T> - Parameters:
-
subscriber- подписчик - Throws:
-
NullPointerException- если подписчик null
submit
public int submit(T item)
Опубликовывает данный элемент всем текущим подписчикам, асинхронно вызывая метод onNext, блокируя не прерывимо, пока ресурсы для любого подписчика недоступны. Этот метод возвращает оценку максимальной задержки (количества элементов, отправленных, но ещё не потреблённых) среди всех текущих подписчиков. Это значение не меньше единицы (учитывая данный отправленный элемент), если есть подписчики, иначе ноль.
Если Executor для этого издателя вызывает исключение RejectedExecutionException (или любое другое RuntimeException или Error) при попытке асинхронного уведомления подписчиков, то это исключение перебрасывается, в этом случае не все подписчики получат этот элемент.
- Parameters:
-
item- публикуемый (не null) элемент - Returns:
- оценка максимальной задержки среди подписчиков
- Throws:
-
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) при попытке асинхронного уведомления подписчиков, или обработчик отбрасывания вызывает исключение при обработке отброшенного элемента, то это исключение перебрасывается.
- Parameters:
-
item- публикуемый (не null) элемент -
onDrop- если не null, обработчик, вызываемый при отбрасывании подписчику, с аргументами подписчик и элемент; если он возвращает true, попытка предложения повторяется (один раз) - Returns:
- если отрицательный, (отрицательное) число отброшенных; в противном случае оценка максимальной задержки
- Throws:
-
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) при попытке асинхронного уведомления подписчиков, или обработчик отбрасывания вызывает исключение при обработке отброшенного элемента, то это исключение перебрасывается.
- Parameters:
-
item- публикуемый (не null) элемент -
timeout- время ожидания ресурсов для любого подписчика перед отказом в единицахunit -
unit-TimeUnitопределяющий, как интерпретировать параметрtimeout -
onDrop- если не null, обработчик, вызываемый при отбрасывании подписчику, с аргументами подписчик и элемент; если он возвращает true, попытка предложения повторяется (один раз) - Returns:
- если отрицательный, (отрицательное) число отброшенных; в противном случае оценка максимальной задержки
- Throws:
-
IllegalStateException- если закрыт -
NullPointerException- если элемент null -
RejectedExecutionException- если вызвано Executorом
close
public void close()
Если не закрыт, отправляет сигналы onComplete текущим подписчикам и запрещает последующие попытки публикации. По возвращении, этот метод НЕ гарантирует, что все подписчики завершили работу.
- Specified by:
-
closeв интерфейсеAutoCloseable
closeExceptionally
public void closeExceptionally(Throwable error)
Если не закрыт, отправляет сигналы onError текущим подписчикам с заданной ошибкой и запрещает последующие попытки публикации. Будущие подписчики также получат данную ошибку. По возвращении, этот метод НЕ гарантирует, что все подписчики завершили работу.
- Parameters:
-
error- аргументonErrorотправленный подписчикам - Throws:
-
NullPointerException- если ошибка null
isClosed
public boolean isClosed()
Возвращает true, если этот издатель не принимает новые сообщения.
- Returns:
- true, если закрыт
getClosedException
public Throwable getClosedException()
Возвращает исключение, связанное с closeExceptionally, или null, если не закрыт или закрыт нормально.
- Returns:
- исключение, или null, если отсутствует
hasSubscribers
public boolean hasSubscribers()
Возвращает true, если у этого издателя есть подписчики.
- Returns:
- true, если у этого издателя есть подписчики
getNumberOfSubscribers
public int getNumberOfSubscribers()
Возвращает количество текущих подписчиков.
- Returns:
- количество текущих подписчиков
getExecutor
public Executor getExecutor()
Возвращает Executor, используемый для асинхронной доставки.
- 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)
Возвращает true, если указанный подписчик сейчас подписан.
- 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)
Обрабатывает все опубликованные элементы, используя заданную функцию Consumer. Возвращает CompletableFuture, который завершается нормально, когда этот издатель сигнализирует onComplete, или завершается с исключением при любой ошибке, или исключение бросается Consumer, или возвращаемый CompletableFuture отменяется, в этом случае дальнейшие элементы не обрабатываются.
- Параметры:
-
consumer- функция, применяемая к каждому элементу onNext - Возвращает:
- CompletableFuture, который завершается успешно, когда издатель сигнализирует об onComplete, и с исключением при любой ошибке или отмене
- Исключения:
-
NullPointerException- если потребитель равен null
© 1993, 2020, 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/11/docs/api/java.base/java/util/concurrent/SubmissionPublisher.html