Класс 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, |
Компонент, который одновременно является 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, которое можно использовать при отсутствии других ограничений. |
Подробное описание методов
defaultBufferSize
public static int defaultBufferSize()
- Примечание по реализации:
- Текущее возвращаемое значение — 256.
- Возвращает:
- значение размера буфера
© 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