Класс 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), и опускает некоторую обработку ошибок, необходимую для полного соответствия спецификации Reactive Streams.
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, которое можно использовать при отсутствии других ограничений. |
Методы, объявленные в классе Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait | Модификатор и тип | Метод | Описание |
|---|---|---|
protected Object |
clone() |
Создаёт и возвращает копию этого объекта. |
boolean |
equals |
Показывает, равен ли этот объект другому объекту. |
protected void |
finalize() |
Устарело, будет удалено: этот элемент API подлежит удалению в будущей версии. Финализация устарела и подлежит удалению в одном из будущих выпусков. |
final Class |
getClass() |
Возвращает класс времени выполнения этого Object. |
int |
hashCode() |
Возвращает хеш-код этого объекта. |
final void |
notify() |
Пробуждает один поток, ожидающий на мониторе этого объекта. |
final void |
notifyAll() |
Пробуждает все потоки, ожидающие на мониторе этого объекта. |
String |
toString() |
Возвращает строковое представление объекта. |
final void |
wait() |
Заставляет текущий поток ожидать пробуждения, обычно в результате уведомления или прерывания. |
final void |
wait |
Заставляет текущий поток ожидать пробуждения, обычно в результате уведомления или прерывания, либо до истечения определённого времени. |
final void |
wait |
Заставляет текущий поток ожидать пробуждения, обычно в результате уведомления или прерывания, либо до истечения определённого времени. |
Подробное описание методов
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.