Класс 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).
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 |
Управление сообщениями, связывающее |
Методы
| Модификатор и тип | Метод | Описание |
|---|---|---|
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