Класс Поток
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)) не выдаются, пока не будут запрошены, но может быть запрошено несколько элементов. Многие реализации подписчика могут организовать это в стиле следующего примера, где размер буфера равен 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, |
Компонент, который одновременно является подписчиком и издателем. |
static interface |
Flow.Publisher<T> |
Производитель элементов (и связанных управляющих сообщений), получаемых подписчиками. |
static interface |
Flow.Subscriber<T> |
Приемник сообщений. |
static interface |
Flow.Subscription |
Управление сообщениями, связывающее Flow.Publisher и Flow.Subscriber. |
Краткое описание методов
| Модификатор и тип | Метод | Описание |
|---|---|---|
static int |
defaultBufferSize() |
Возвращает значение по умолчанию для буферизации издателя или подписчика, которое может быть использовано в отсутствие других ограничений. |
Подробное описание методов
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://download.java.net/java/early_access/jdk24/docs/api/java.base/java/util/concurrent/Flow.html