Spec-Zone.ru › OpenJDK 21

Класс 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)) не выдаются, пока не запрошены, но может быть запрошено несколько элементов. Многие реализации подписчика могут организовать это в стиле следующего примера, где размер буфера 1 однозначно определяет шаги, а большие размеры обычно позволяют более эффективную перекрывающуюся обработку с меньшим объёмом связи; например, со значением 64, это поддерживает общее количество ожидающих запросов между 32 и 64. Поскольку вызовы методов подписчика для данного Flow.Subscription строго упорядочены, этим методам нет необходимости использовать блокировки или volatile, если подписчик не поддерживает несколько подписок (в этом случае лучше определить несколько подписчиков, каждый со своей подписью).

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

 
 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) { ... }
 }
Since:
9

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

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

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

Модификатор и тип Метод Описание
static int defaultBufferSize()
Возвращает значение по умолчанию для буферизации издателя или подписчика, которое может использоваться при отсутствии других ограничений.

Методы, унаследованные от класса java.lang.Object

clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait

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

defaultBufferSize

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

© 1993, 2023, 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/21/docs/api/java.base/java/util/concurrent/Flow.html

Spec-Zone.ru

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