PHpullh
학습 라이브러리/Java/리액티브 스트림 연산자

JAVA · 비동기/동시성

리액티브 스트림 연산자

SubmissionPublisherProcessor를 사용한 데이터 변환 파이프라인입니다.

비동기/동시성고급Processorreactive파이프라인변환

핵심 설명

SubmissionPublisherProcessor를 사용한 데이터 변환 파이프라인입니다.

Java code

import java.util.concurrent.*;
import java.util.concurrent.Flow.*;
import java.util.function.Function;

public class ReactiveProcessor {
    // Processor — Publisher + Subscriber (변환기)
    static class TransformProcessor<T, R>
            extends SubmissionPublisher<R>
            implements Processor<T, R> {

        private final Function<T, R> transform;
        private Subscription subscription;

        TransformProcessor(Function<T, R> transform) {
            this.transform = transform;
        }

        @Override
        public void onSubscribe(Subscription s) {
            this.subscription = s;
            s.request(1);
        }

        @Override
        public void onNext(T item) {
            submit(transform.apply(item)); // 변환 후 전달
            subscription.request(1);
        }

        @Override
        public void onError(Throwable t) { closeExceptionally(t); }

        @Override
        public void onComplete() { close(); }
    }

    public static void main(String[] args) throws Exception {
        var publisher = new SubmissionPublisher<String>();
        var upper = new TransformProcessor<String, String>(String::toUpperCase);
        var lengthProc = new TransformProcessor<String, Integer>(String::length);

        // 파이프라인: publisher -> upper -> lengthProc -> subscriber
        publisher.subscribe(upper);
        upper.subscribe(lengthProc);
        lengthProc.subscribe(new Flow.Subscriber<>() {
            public void onSubscribe(Subscription s) { s.request(Long.MAX_VALUE); }
            public void onNext(Integer item) { System.out.println("길이: " + item); }
            public void onError(Throwable t) {}
            public void onComplete() { System.out.println("완료!"); }
        });

        java.util.List.of("hello", "world", "java").forEach(publisher::submit);
        publisher.close();
        Thread.sleep(500);
    }
}

학습 팁

ProcessorPublisherSubscriber를 모두 구현합니다. 데이터 변환, 필터링, 버퍼링 등의 중간 처리에 사용합니다.

주의할 점

Processor에서 request(1)을 호출하지 않으면 다음 데이터를 받을 수 없습니다. onNext에서 반드시 추가 요청을 하세요.

자주 묻는 질문

리액티브 스트림 연산자란 무엇인가요?

SubmissionPublisher 와 Processor 를 사용한 데이터 변환 파이프라인입니다.

리액티브 스트림 연산자 학습 시 주의할 점은 무엇인가요?

Processor 에서 request(1) 을 호출하지 않으면 다음 데이터를 받을 수 없습니다. onNext 에서 반드시 추가 요청을 하세요.