본문으로 건너뛰기
김신건의 로그

[Java] Phaser

· 수정 · 📖 약 2분 · 799자/단어 #java #concurrent #synchronizer #barrier #dynamic
Phaser, java.util.concurrent.Phaser, 페이저, 동적 배리어

정의

java.util.concurrent.PhaserCyclicBarrier 의 확장. 참여자 수가 동적으로 변할 수 있는 단계별 동기화. JDK 1.7 도입.

시각화

핵심 메서드

Phaser phaser = new Phaser(N);

phaser.arriveAndAwaitAdvance();    // 도착 + 다른 참여자 대기
phaser.arrive();                    // 도착만 (대기 안 함)
phaser.awaitAdvance(phase);         // 특정 phase 까지 대기
phaser.register();                  // 참여자 추가
phaser.arriveAndDeregister();       // 도착 + 자신 제거

phaser.getPhase();                  // 현재 phase 번호
phaser.getRegisteredParties();      // 등록된 참여자 수
phaser.getArrivedParties();         // 이미 도착한 참여자 수
phaser.getUnarrivedParties();       // 아직 도착 안 한 참여자 수
phaser.isTerminated();              // phaser 종료 여부

동적 참여 예

Phaser phaser = new Phaser(1);   // 시작 시 main 만 등록

for (int i = 0; i < workers; i++) {
    phaser.register();           // worker 추가
    new Thread(() -> {
        while (running) {
            doPhaseWork();
            phaser.arriveAndAwaitAdvance();   // 다른 worker 대기
        }
        phaser.arriveAndDeregister();         // 종료 시 deregister
    }).start();
}

phaser.arriveAndDeregister();   // main 도 제거

worker 가 도중에 추가/종료될 수 있는 경우 CyclicBarrier 로는 표현 어렵다.

onAdvance 콜백

Phaser phaser = new Phaser() {
    @Override
    protected boolean onAdvance(int phase, int registered) {
        System.out.println("phase " + phase + " done");
        return phase >= 10 || registered == 0;   // true 면 종료
    }
};

각 phase 가 끝날 때 호출. true 반환하면 phaser 종료.

CountDownLatch / CyclicBarrier / Phaser 비교

기능CountDownLatchCyclicBarrierPhaser
재사용
참여자 동적 변경
한 명만 도착 (대기 안 함)✓ (arrive)
콜백barrier actiononAdvance
트리 구조 (계층)✓ (parent phaser)

가장 유연하지만 그만큼 복잡. 단순 N → 0 은 CountDownLatch, 고정 N 반복은 CyclicBarrier 가 더 간단.

다단계 배리어 흐름

sequenceDiagram
    participant Main
    participant W1 as Worker 1
    participant W2 as Worker 2
    participant P as Phaser

    Main->>P: new Phaser(1)
    Main->>P: register()
    Main->>W1: Thread.start()
    Main->>P: register()
    Main->>W2: Thread.start()
    Main->>P: arriveAndDeregister()

    W1->>P: arriveAndAwaitAdvance() Phase 0
    W2->>P: arriveAndAwaitAdvance() Phase 0
    Note over P: Phase 0 완료, onAdvance() 호출
    P-->>W1: 재개 (Phase 1)
    P-->>W2: 재개 (Phase 1)

    W1->>P: arriveAndDeregister()
    W2->>P: arriveAndAwaitAdvance() Phase 1
    Note over P: W1 이탈, parties=1, Phase 1 완료
    P-->>W2: 재개 (Phase 2)
    W2->>P: arriveAndDeregister()
    Note over P: 모든 참여자 탈퇴, 종료

각 phase 마다 도착해야 하는 참여자 수가 달라질 수 있음을 보여준다.

내부 구조 (분산 카운터)

Phaser 는 단일 long 필드 (state) 에 여러 값을 비트 패킹한다.

state = [unarrived (16b)] [parties (16b)] [phase (31b)] [terminated (1b)]
  • unarrived: 아직 arrive 하지 않은 참여자 수
  • parties: 등록된 총 참여자 수
  • phase: 현재 phase 번호 (0 ~ Integer.MAX_VALUE 순환)
  • terminated: phaser 종료 여부

모든 상태 전이는 Unsafe.compareAndSwapLong (CAS) 으로 이루어진다. CyclicBarrier 처럼 전역 락이 없고, arrive 는 lock-free 다. 단, 대기 중인 스레드의 park/unpark 는 내부 큐로 관리.

계층적 Phaser (parent 지정 시) 는 마지막 unarrived 가 0 이 되면 자동으로 parent 에 arrive 를 위임한다.

forceTermination

phaser.forceTermination();

진행 중인 모든 phase 를 중단하고 phaser 를 즉시 종료 상태로 만든다. awaitAdvance() 로 대기 중이던 스레드는 음수 phase (terminated marker) 를 받으며 깨어난다.

// 취소 신호 처리 패턴
void cancelAll(Phaser phaser) {
    phaser.forceTermination();
    // 이후 isTerminated() == true
    // getPhase() 반환값이 음수가 됨
}

// 종료 감지
int phase = phaser.arriveAndAwaitAdvance();
if (phase < 0) {
    // phaser 가 terminated, 정리 작업
}

계층적 Phaser (parent 지정)

대규모 병렬 작업에서 단일 Phaser 의 CAS 경합이 심해질 때, 트리 구조의 Phaser 로 분산시킬 수 있다.

// 1024개 스레드: 32개 그룹, 각 32 스레드
Phaser root = new Phaser();

List<Phaser> children = new ArrayList<>();
for (int g = 0; g < 32; g++) {
    Phaser child = new Phaser(root, 32);  // parent = root, parties = 32
    children.add(child);
}

// 각 그룹의 스레드는 child phaser 를 사용
// 그룹 내 32명이 모두 도착하면 root 에 arrive 전달
// root 에서 32개 그룹 전체가 도착해야 phase 전환

이 방식으로 O(스레드 수) CAS 경합을 O(log 스레드 수) 수준으로 줄일 수 있다. ForkJoinPool 의 병렬 처리와 유사한 트리 분할.

실전 예시: 다단계 ETL 파이프라인

// Java 17+
record Worker(int id, List<String> data) implements Runnable {

    @Override
    public void run() {
        Phaser phaser = PipelineContext.phaser();

        // Phase 0: 데이터 추출
        var extracted = extract(data);
        phaser.arriveAndAwaitAdvance();

        // Phase 1: 변환 (모든 worker 의 추출이 끝난 뒤)
        var transformed = transform(extracted);
        phaser.arriveAndAwaitAdvance();

        // Phase 2: 적재
        load(transformed);
        phaser.arriveAndDeregister();  // 완료, 참여 종료
    }
}

// 메인
Phaser etlPhaser = new Phaser(1) {
    @Override
    protected boolean onAdvance(int phase, int registered) {
        System.out.printf("[ETL] Phase %d 완료, 참여자=%d%n", phase, registered);
        return registered == 0;   // 모두 deregister 되면 종료
    }
};

for (int i = 0; i < WORKER_COUNT; i++) {
    etlPhaser.register();
    new Thread(new Worker(i, chunks.get(i))).start();
}
etlPhaser.arriveAndDeregister();   // main thread 제거

JMM 보장

  • arriveAndAwaitAdvance() 반환 이후에는 이전 phase 에서 다른 모든 스레드가 쓴 값이 visible
  • arrive 시 CAS (내부 Unsafe.compareAndSwapLong) 가 happens-before fence 역할
  • phase 경계를 넘으면 공유 상태의 안전한 가시성이 보장됨

성능 특성

항목특성
arrive (lock-free)CAS 기반, 경합 없을 때 매우 빠름
대기 (awaitAdvance)park 기반, CPU 낭비 없음
확장성단일 Phaser: 수백 스레드까지 양호, 이상이면 계층 구조 권장
state 필드단일 long 비트 패킹, 원자 갱신

함정

1. arrive 와 await 는 별개

arrive() 는 도착 신호만 보내고 대기하지 않는다. awaitAdvance(phaser.getPhase()) 를 별도로 호출해야 다음 phase 시작을 기다릴 수 있다.

2. unarrived 가 0 이 돼야 phase 전환

Phaser phaser = new Phaser(2);
phaser.arrive();        // parties=2, unarrived=1, phase=0
// phase 는 아직 0! 두 번째 arrive 가 있어야 phase 1 로 전환
phaser.arrive();        // 이제 phase=1

3. onAdvance 에서 예외 발생 시 종료

onAdvance 에서 unchecked exception 이 발생하면 phaser 가 즉시 terminated 된다. 예외 처리를 onAdvance 내부에서 완결해야 한다.

4. register 없이 arrive 시 IllegalStateException

Phaser phaser = new Phaser(0);  // 참여자 0
phaser.arrive();                 // IllegalStateException: parties=0

arrive 전에 반드시 register() 또는 생성자 인자로 parties 를 등록해야 한다.

관련 위키

이 글의 용어 (4개)
[Java] CountDownLatchjava
정의 는 한 스레드 또는 여러 스레드가 다른 스레드들의 작업 완료를 기다리는 일회용 synchronizer. 으로 초기화하고, 각 worker 가 작업 후 호출. waiter 는…
[Java] CyclicBarrierjava
정의 는 N 개의 스레드가 모두 도착할 때까지 서로 기다리는 synchronizer. 모두 도착하면 동시에 출발, 다시 사용 가능 (cyclic). 시각화 핵심 메서드 가장 흔한…
[Java] Semaphorejava
정의 는 N 개의 permit (허가) 을 발급/회수 하는 synchronizer. permit 이 있으면 진행, 없으면 또는 즉시 실패. 자원 풀, rate limiting, …
Blocking (블로킹)concurrency
정의 Blocking 은 호출이 결과가 준비될 때까지 호출 스레드를 멈춰두는 실행 방식. 멈춰 있는 동안 그 스레드는 다른 일을 할 수 없다. OS 가 스레드를 sleep 상태로…

💬 댓글

사이트 검색 / 명령어

검색

스크롤 = 확대/축소 · 드래그 = 이동 · 0 = 원래 크기 · ESC = 닫기