Spec-Zone.ru › OpenJDK 27

Класс 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
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(Executor executor, int maxBufferCapacity)
Создаёт новый SubmissionPublisher, использующий указанный Executor для асинхронной доставки подписчикам, с указанным максимальным размером буфера для каждого подписчика и без обработчика исключений Subscriber в методе onNext.
SubmissionPublisher(Executor executor, int maxBufferCapacity, BiConsumer<? super Flow.Subscriber<? super T>, ? super Throwable> handler)
Создаёт новый SubmissionPublisher, использующий указанный Executor для асинхронной доставки подписчикам, с указанным максимальным размером буфера для каждого подписчика и, если он не равен null, указанным обработчиком, вызываемым при выбрасывании любым Subscriber исключения в методе onNext.

Краткое описание методов

Модификатор и тип Метод Описание
void close()
Если издатель ещё не закрыт, отправляет текущим подписчикам сигналы onComplete и запрещает последующие попытки публикации.
void closeExceptionally(Throwable error)
Если издатель ещё не закрыт, отправляет текущим подписчикам сигналы onError с указанной ошибкой и запрещает последующие попытки публикации.
CompletableFuture<Void> consume(Consumer<? super T> consumer)
Обрабатывает все опубликованные элементы с помощью указанной функции 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(Flow.Subscriber<? super T> subscriber)
Возвращает true, если указанный Subscriber в данный момент подписан.
int offer(T item, long timeout, TimeUnit unit, BiPredicate<Flow.Subscriber<? super T>, ? super T> onDrop)
Если возможно, публикует указанный элемент для каждого текущего подписчика, асинхронно вызывая его метод onNext. Блокируется, пока ресурсы для какой-либо подписки недоступны, до истечения указанного времени ожидания или прерывания потока вызывающего кода; после этого вызывается указанный обработчик (если он не равен null), и, если он возвращает true, выполняется одна повторная попытка.
int offer(T item, BiPredicate<Flow.Subscriber<? super T>, ? super T> onDrop)
Если возможно, публикует указанный элемент для каждого текущего подписчика, асинхронно вызывая его метод onNext.
int submit(T item)
Публикует указанный элемент для каждого текущего подписчика, асинхронно вызывая его метод onNext, и без возможности прерывания блокируется, пока ресурсы для любого подписчика недоступны.
void subscribe(Flow.Subscriber<? super T> subscriber)
Добавляет указанного Subscriber, если он ещё не подписан.

Методы, объявленные в классе Object

clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
Модификатор и тип Метод Описание
protected Object clone()
Создаёт и возвращает копию этого объекта.
boolean equals(Object obj)
Указывает, равен ли этот объект какому-либо другому объекту.
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(long timeoutMillis)
Заставляет текущий поток ожидать пробуждения, обычно вследствие вызова notify или interrupt, либо до истечения заданного промежутка реального времени.
final void wait(long timeoutMillis, int nanos)
Заставляет текущий поток ожидать пробуждения, обычно вследствие вызова notify или interrupt, либо до истечения заданного промежутка реального времени.

Подробное описание конструкторов

SubmissionPublisher

public SubmissionPublisher(Executor executor, int maxBufferCapacity, BiConsumer<? super Flow.Subscriber<? super T>, ? super Throwable> handler)
Создаёт новый SubmissionPublisher, использующий указанный Executor для асинхронной доставки подписчикам, с указанным максимальным размером буфера для каждого подписчика и, если он не равен null, указанным обработчиком, вызываемым при выбрасывании любым Subscriber исключения в методе onNext.
Параметры:
executor — исполнитель для асинхронной доставки, поддерживающий создание как минимум одного независимого потока
maxBufferCapacity — максимальная ёмкость буфера каждого подписчика (фактическая ёмкость может быть округлена вверх до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым этой реализацией; метод getMaxBufferCapacity() возвращает фактическое значение)
handler — если не равен null, процедура, вызываемая при исключении, выброшенном в методе onNext
Исключения:
NullPointerException — если executor равен null
IllegalArgumentException — если maxBufferCapacity не положителен

SubmissionPublisher

public SubmissionPublisher(Executor executor, int maxBufferCapacity)
Создаёт новый SubmissionPublisher, использующий указанный Executor для асинхронной доставки подписчикам, с указанным максимальным размером буфера для каждого подписчика и без обработчика исключений Subscriber в методе onNext.
Параметры:
executor — исполнитель для асинхронной доставки, поддерживающий создание как минимум одного независимого потока
maxBufferCapacity — максимальная ёмкость буфера каждого подписчика (фактическая ёмкость может быть округлена вверх до ближайшей степени двойки и/или ограничена наибольшим значением, поддерживаемым этой реализацией; метод getMaxBufferCapacity() возвращает фактическое значение)
Исключения:
NullPointerException — если executor равен null
IllegalArgumentException — если maxBufferCapacity не положителен

SubmissionPublisher

public SubmissionPublisher()
Создаёт новый SubmissionPublisher, использующий ForkJoinPool.commonPool() для асинхронной доставки подписчикам, с максимальной ёмкостью буфера Flow.defaultBufferSize() и без обработчика исключений Subscriber в методе onNext.

Подробное описание методов

subscribe

public void subscribe(Flow.Subscriber<? super T> subscriber)
Добавляет указанного 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, если этот издатель не принимает новые элементы.
Возвращает:
true, если издатель закрыт

getClosedException

public Throwable getClosedException()
Возвращает исключение, связанное с closeExceptionally, или null, если издатель не закрыт либо закрыт штатно.
Возвращает:
исключение или null, если его нет

hasSubscribers

public boolean hasSubscribers()
Возвращает true, если у этого издателя есть подписчики.
Возвращает:
true, если у этого издателя есть подписчики

getNumberOfSubscribers

public int getNumberOfSubscribers()
Возвращает число текущих подписчиков.
Возвращает:
число текущих подписчиков

getExecutor

public Executor getExecutor()
Возвращает Executor, используемый для асинхронной доставки.
Возвращает:
Executor, используемый для асинхронной доставки

getMaxBufferCapacity

public int getMaxBufferCapacity()
Возвращает максимальную ёмкость буфера для каждого подписчика.
Возвращает:
максимальную ёмкость буфера для каждого подписчика

getSubscribers

public List<Flow.Subscriber<? super T>> getSubscribers()
Возвращает список текущих подписчиков для мониторинга и отслеживания, а не для вызова методов Flow.Subscriber у подписчиков.
Возвращает:
список текущих подписчиков

isSubscribed

public boolean isSubscribed(Flow.Subscriber<? super T> subscriber)
Возвращает true, если указанный Subscriber в данный момент подписан.
Параметры:
subscriber — подписчик
Возвращает:
true, если подписка активна
Исключения:
NullPointerException — если subscriber равен null

estimateMinimumDemand

public long estimateMinimumDemand()
Возвращает оценку минимального числа элементов, запрошенных (посредством request), но ещё не произведённых всеми текущими подписчиками.
Возвращает:
оценку или ноль, если подписчиков нет

estimateMaximumLag

public int estimateMaximumLag()
Возвращает оценку максимального числа элементов, произведённых, но ещё не потреблённых всеми текущими подписчиками.
Возвращает:
оценку

consume

public CompletableFuture<Void> consume(Consumer<? super T> consumer)
Обрабатывает все опубликованные элементы с помощью указанной функции Consumer. Возвращает CompletableFuture, который завершается штатно, когда этот издатель отправляет сигнал onComplete, или завершается исключением при любой ошибке, при выбрасывании исключения Consumer либо при отмене возвращённого CompletableFuture; в последнем случае дальнейшая обработка элементов не выполняется.
Параметры:
consumer — функция, применяемая к каждому элементу onNext
Возвращает:
CompletableFuture, который завершается штатно при отправке издателем сигнала onComplete и завершается исключением при любой ошибке или отмене
Исключения:
NullPointerException — если consumer равен null

Сообщить об ошибке или предложить улучшение
Дополнительную справочную информацию по API и документацию для разработчиков см. в документации Java SE, содержащей более подробные описания для разработчиков, включая концептуальные обзоры, определения терминов, обходные решения и рабочие примеры кода. Другие версии.
Java является товарным знаком или зарегистрированным товарным знаком Oracle и/или ее дочерних компаний в США и других странах.
Авторские права © 1993, 2026, Oracle и/или ее дочерние компании, 500 Oracle Parkway, Redwood Shores, CA 94065, США.
Все права защищены. Использование регулируется условиями лицензии и политикой распространения документации.

© 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.

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API