Класс 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).

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) { ... }
 }
С момента:
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()

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

Примечание к реализации:
Текущее возвращаемое значение равно 256.
Возвращает:
значение размера буфера

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

Spec-Zone .ru
спецификации, руководства, описания, API