Spec-Zone.ru › OpenJDK 27

Класс Flow

java.lang.Object
java.util.concurrent.Flow
public final class Flow extends Object
Связанные интерфейсы и статические методы для создания компонентов с управлением потоком, в которых Publishers создают элементы, потребляемые одним или несколькими Subscribers, каждый из которых управляется объектом Subscription.

Эти интерфейсы соответствуют спецификации reactive-streams. Они применяются как в параллельных, так и в распределённых асинхронных средах: все (семь) методов определены в стиле сообщений «однонаправленного» обмена void. Взаимодействие основано на простой форме управления потоком (методе Flow.Subscription.request(long)), который можно использовать, чтобы избежать проблем с управлением ресурсами, которые иначе могут возникнуть в системах на основе «push»-модели.

Примеры. Flow.Publisher обычно определяет собственную реализацию Flow.Subscription, создаёт её в методе subscribe и передаёт вызывающему коду Flow.Subscriber. Он асинхронно публикует элементы для подписчика, обычно с помощью Executor. Например, ниже приведён очень простой издатель, который выдаёт (по запросу) единственный элемент TRUE одному подписчику. Поскольку подписчик получает только один элемент, этот класс не использует буферизацию и управление порядком, необходимые в большинстве реализаций (например, SubmissionPublisher), и опускает некоторую обработку ошибок, необходимую для полного соответствия спецификации Reactive Streams.

class OneShotPublisher implements Publisher<Boolean> {
  private final ExecutorService executor = ForkJoinPool.commonPool(); // daemon-based
  private boolean subscribed; // true after first subscribe
  public synchronized void subscribe(Subscriber<? super Boolean> subscriber) {
    if (subscribed)
      subscriber.onError(new IllegalStateException()); // only one allowed
    else {
      subscribed = true;
      subscriber.onSubscribe(new OneShotSubscription(subscriber, executor));
    }
  }
  static class OneShotSubscription implements Subscription {
    private final Subscriber<? super Boolean> subscriber;
    private final ExecutorService executor;
    private Future<?> future; // to allow cancellation
    private boolean completed;
    OneShotSubscription(Subscriber<? super Boolean> subscriber,
                        ExecutorService executor) {
      this.subscriber = subscriber;
      this.executor = executor;
    }
    public synchronized void request(long n) {
      if (!completed) {
        completed = true;
        if (n <= 0) {
          IllegalArgumentException ex = new IllegalArgumentException();
          executor.execute(() -> subscriber.onError(ex));
        } else {
          future = executor.submit(() -> {
            subscriber.onNext(Boolean.TRUE);
            subscriber.onComplete();
          });
        }
      }
    }
    public synchronized void cancel() {
      completed = true;
      if (future != null) future.cancel(false);
    }
  }
}

Flow.Subscriber организует запрос и обработку элементов. Элементы (вызовы Flow.Subscriber.onNext(T)) не отправляются, если их не запросили, но можно запросить несколько элементов. Многие реализации Subscriber могут организовать это так, как показано в следующем примере: размер буфера 1 обеспечивает пошаговую обработку, а большие размеры обычно позволяют эффективнее перекрывать обработку и уменьшить количество взаимодействий; например, значение 64 поддерживает общее число невыполненных запросов в диапазоне от 32 до 64. Поскольку вызовы методов Subscriber для данного объекта Flow.Subscription строго упорядочены, этим методам не нужно использовать блокировки или volatile-переменные, если только Subscriber не управляет несколькими подписками (в этом случае лучше определить несколько Subscriber, каждый со своей подпиской).

class SampleSubscriber<T> implements Subscriber<T> {
  final Consumer<? super T> consumer;
  Subscription subscription;
  final long bufferSize;
  long count;
  SampleSubscriber(long bufferSize, Consumer<? super T> consumer) {
    this.bufferSize = bufferSize;
    this.consumer = consumer;
  }
  public void onSubscribe(Subscription subscription) {
    long initialRequestSize = bufferSize;
    count = bufferSize - bufferSize / 2; // re-request when half consumed
    (this.subscription = subscription).request(initialRequestSize);
  }
  public void onNext(T item) {
    if (--count <= 0)
      subscription.request(count = bufferSize - bufferSize / 2);
    consumer.accept(item);
  }
  public void onError(Throwable ex) { ex.printStackTrace(); }
  public void onComplete() {}
}

Значение по умолчанию defaultBufferSize() может стать полезной отправной точкой при выборе размеров запросов и буферов компонентов Flow с учётом ожидаемой скорости, ресурсов и сценариев использования. Если же управление потоком никогда не требуется, подписчик может изначально запросить практически неограниченное количество элементов, как показано ниже:

class UnboundedSubscriber<T> implements Subscriber<T> {
  public void onSubscribe(Subscription subscription) {
    subscription.request(Long.MAX_VALUE); // effectively unbounded
  }
  public void onNext(T item) { use(item); }
  public void onError(Throwable ex) { ex.printStackTrace(); }
  public void onComplete() {}
  void use(T item) { ... }
}
Начиная с:
9

Краткое описание вложенных классов

Модификатор и тип Класс Описание
static interface  Flow.Processor<T,R>
Компонент, выполняющий функции как Subscriber, так и Publisher.
static interface  Flow.Publisher<T>
Источник элементов (и связанных управляющих сообщений), получаемых подписчиками.
static interface  Flow.Subscriber<T>
Получатель сообщений.
static interface  Flow.Subscription
Управляющее сообщение, связывающее Flow.Publisher и Flow.Subscriber.

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

Модификатор и тип Метод Описание
static int defaultBufferSize()
Возвращает значение по умолчанию для буферизации Publisher или 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()
Заставляет текущий поток ожидать пробуждения, обычно в результате уведомления или прерывания.
final void wait(long timeoutMillis)
Заставляет текущий поток ожидать пробуждения, обычно в результате уведомления или прерывания, либо до истечения определённого времени.
final void wait(long timeoutMillis, int nanos)
Заставляет текущий поток ожидать пробуждения, обычно в результате уведомления или прерывания, либо до истечения определённого времени.

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

defaultBufferSize

public static int defaultBufferSize()
Возвращает значение по умолчанию для буферизации Publisher или Subscriber, которое можно использовать при отсутствии других ограничений.
Примечание по реализации:
Текущее возвращаемое значение — 256.
Возвращает:
значение размера буфера

Сообщить об ошибке или предложить улучшение
Дополнительную справочную информацию по API и документацию для разработчиков см. в документации Java SE, содержащей более подробные описания для разработчиков, в том числе концептуальные обзоры, определения терминов, обходные решения и рабочие примеры кода. Другие версии.
Java является товарным знаком или зарегистрированным товарным знаком Oracle и/или её аффилированных лиц в США и других странах.
Авторское право © 1993, 2026, 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.

Spec-Zone.ru

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