Что такое "backpressure" в реактивных системах на Java, и как его можно обработать?
Ответ ⬇️
"Backpressure" возникает, когда источник данных производит элементы быстрее, чем потребитель может их обработать. В Java для управления этим в реактивных системах используется интерфейс Flow.Subscriber, который позволяет контролировать поток данных через запрос определенного количества элементов с помощью метода request.
🗣 Пример:
import java.util.concurrent.Flow;
public class BackpressureExample {
public static void main(String[] args) {
MyPublisher publisher = new MyPublisher();
MySubscriber subscriber = new MySubscriber();
publisher.subscribe(subscriber);
}
}
class MyPublisher implements Flow.Publisher<Integer> {
@Override
public void subscribe(Flow.Subscriber<? super Integer> subscriber) {
subscriber.onSubscribe(new Flow.Subscription() {
@Override
public void request(long n) {
for (int i = 1; i <= n; i++) {
subscriber.onNext(i);
}
subscriber.onComplete();
}
@Override
public void cancel() {
}
});
}
}
class MySubscriber implements Flow.Subscriber<Integer> {
private Flow.Subscription subscription;
@Override
public void onSubscribe(Flow.Subscription subscription) {
this.subscription = subscription;
subscription.request(5); // Запрашиваем 5 элементов
}
@Override
public void onNext(Integer item) {
System.out.println("Получен элемент: " + item);
}
@Override
public void onError(Throwable throwable) {
System.err.println("Ошибка: " + throwable.getMessage());
}
@Override
public void onComplete() {
System.out.println("Все элементы получены.");
}
}
Java Learning 👩💻