JAVA · 비동기/동시성
리액티브 스트림과 배압
리액티브 스트림의 배압(backpressure) 메커니즘을 이해합니다.
비동기/동시성고급reactivebackpressure배압SubmissionPublisher
핵심 설명
리액티브 스트림의 배압(backpressure) 메커니즘을 이해합니다.
Java code
import java.util.concurrent.*;
import java.util.concurrent.Flow.*;
public class BackpressureDemo {
// 배압 지원 Subscriber
static class ControlledSubscriber implements Subscriber<Integer> {
private Subscription subscription;
private int received = 0;
private final int batchSize;
ControlledSubscriber(int batchSize) { this.batchSize = batchSize; }
@Override
public void onSubscribe(Subscription s) {
this.subscription = s;
s.request(batchSize); // 초기 요청량
}
@Override
public void onNext(Integer item) {
System.out.println("처리: " + item);
received++;
if (received % batchSize == 0) {
// 배치 처리 완료 후 다음 배치 요청
System.out.println("-- 다음 배치 요청 --");
subscription.request(batchSize);
}
}
@Override
public void onError(Throwable t) {
System.err.println("에러: " + t);
}
@Override
public void onComplete() {
System.out.println("스트림 완료! 총 " + received + "개");
}
}
public static void main(String[] args) throws Exception {
// SubmissionPublisher — JDK 내장 Publisher
try (var publisher = new SubmissionPublisher<Integer>()) {
publisher.subscribe(new ControlledSubscriber(3));
for (int i = 1; i <= 10; i++) {
publisher.submit(i);
}
} // close() -> onComplete() 호출
Thread.sleep(1000); // 비동기 처리 대기
}
}학습 팁
SubmissionPublisher는 JDK에 포함된 유일한 Publisher 구현입니다. 프로덕션에서는 Project Reactor를 사용하세요.
주의할 점
request(Long.MAX_VALUE)는 사실상 배압 없음입니다. 처리 능력에 맞는 요청량을 설정하지 않으면 메모리 폭주가 발생합니다.
자주 묻는 질문
리액티브 스트림과 배압란 무엇인가요?
리액티브 스트림의 배압(backpressure) 메커니즘을 이해합니다.
리액티브 스트림과 배압 학습 시 주의할 점은 무엇인가요?
request(Long.MAX_VALUE) 는 사실상 배압 없음입니다. 처리 능력에 맞는 요청량을 설정하지 않으면 메모리 폭주가 발생합니다.
Continue Learning
Java 학습을 이어가세요
총 200개의 독립 HTML 학습 문서 중 하나입니다. 각 문서는 고유 URL과 canonical 메타데이터를 가집니다.