Класс Поток
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, |
Компонент, который действует как подписчик и издатель. |
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, 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