Spec-Zone.ru › OpenJDK 17

Класс Поток

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

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

Примеры. 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, 2021, 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/17/docs/api/java.base/java/util/concurrent/Flow.html

Spec-Zone.ru

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