Spec-Zone.ru › OpenJDK 25

Класс Flow

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

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

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

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

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

defaultBufferSize

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

Сообщить об ошибке или предложить улучшение
Дополнительную справочную информацию по 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/Flow.html

Spec-Zone.ru

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