Spec-Zone.ru › OpenJDK 25

Класс 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, который асинхронно передает отправленные (ненулевые) элементы текущим подписчикам до закрытия. Каждый текущий подписчик получает вновь отправленные элементы в том же порядке, если не происходит отбрасывание элементов или исключений. Использование 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(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

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

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 этого 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, если этот издатель не принимает отправку элементов.
Возвращает:
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, 2025, Oracle и/или ее филиалы, 500 Oracle Parkway, Redwood Shores, CA 94065 USA.
Все права защищены. Использование регулируется условиями лицензии и политикой распространения документации.

© 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

Spec-Zone.ru

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