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

[Java] DelayQueue

· 수정 · 📖 약 2분 · 688자/단어 #java #concurrent #queue #blocking #delay #scheduling
DelayQueue, java.util.concurrent.DelayQueue, 지연 큐

정의

java.util.concurrent.DelayQueue<E extends Delayed>각 원소가 지정된 delay 가 만료된 후에만 꺼낼 수 있는 BlockingQueue. 우선순위는 만료 시각이 빠른 것부터.

내부는 PriorityQueue (min-heap). 단일 ReentrantLock 으로 thread-safe.

사용 상황

  • TTL 캐시: 만료된 항목을 별도 스레드가 주기적으로 정리
  • 재시도 큐 (retry queue): 실패한 요청을 일정 시간 뒤 재시도
  • 작업 스케줄러: ScheduledThreadPoolExecutor 와 유사한 단순 구현
  • 세션 만료: 로그인 세션, 인증 토큰 만료 처리

시각화: 시간 기반 꺼내기

flowchart LR
    subgraph queue["DelayQueue (min-heap by expiresAt)"]
        A["Task A\nexpiresAt: T+1s"]
        B["Task B\nexpiresAt: T+3s"]
        C["Task C\nexpiresAt: T+5s"]
        A --> B --> C
    end
    TAKE["take()"] -->|"T+1s 도달 시 반환"| A
    TAKE -->|"T+3s 도달 시 반환"| B
    TAKE -->|"T+5s 도달 시 반환"| C

Delayed 인터페이스

원소는 Delayed 인터페이스를 구현해야 한다.

public interface Delayed extends Comparable<Delayed> {
    long getDelay(TimeUnit unit);
}

getDelay남은 시간 을 반환. 0 이하가 되면 take 가능. compareTo 는 힙 정렬에 사용.

// Java 17+: record 기반 구현
record DelayedTask(String name, long readyAt, Runnable job)
        implements Delayed {

    DelayedTask(String name, long delayMs, Runnable job) {
        this(name, System.currentTimeMillis() + delayMs, job);
    }

    @Override
    public long getDelay(TimeUnit unit) {
        return unit.convert(readyAt - System.currentTimeMillis(), TimeUnit.MILLISECONDS);
    }

    @Override
    public int compareTo(Delayed o) {
        return Long.compare(this.readyAt, ((DelayedTask) o).readyAt());
    }
}

동작

DelayQueue<DelayedTask> q = new DelayQueue<>();
q.put(new DelayedTask("slow", 5000, () -> System.out.println("5초 뒤")));
q.put(new DelayedTask("fast", 1000, () -> System.out.println("1초 뒤")));

// take() 는 만료 순으로 반환: 1s -> 5s
DelayedTask t = q.take();   // 1초 대기 후 "fast" 반환
t.job().run();
t = q.take();               // 추가 4초 대기 후 "slow" 반환
t.job().run();

take() 는 큐의 head 원소의 getDelay 가 0 이하가 될 때까지 block.

복잡도

내부는 PriorityQueue (heap).

작업시간
put, offerO(log n)
take, pollO(log n)
peekO(1)
sizeO(1)

실전 코드: TTL 캐시

public class TtlCache<K, V> {
    private final Map<K, V> store = new ConcurrentHashMap<>();
    private final DelayQueue<ExpiryKey<K>> expirations = new DelayQueue<>();

    // 만료 키 추적
    record ExpiryKey<K>(K key, long expiresAt) implements Delayed {
        ExpiryKey(K key, long ttlMs) {
            this(key, System.currentTimeMillis() + ttlMs);
        }

        @Override
        public long getDelay(TimeUnit unit) {
            return unit.convert(expiresAt - System.currentTimeMillis(), MILLISECONDS);
        }

        @Override
        public int compareTo(Delayed o) {
            return Long.compare(expiresAt, ((ExpiryKey<?>) o).expiresAt());
        }
    }

    public TtlCache() {
        // 만료 스캔 스레드 (데몬)
        Thread sweeper = Thread.ofVirtual().start(() -> {
            while (!Thread.currentThread().isInterrupted()) {
                try {
                    ExpiryKey<K> expired = expirations.take();
                    store.remove(expired.key());
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                }
            }
        });
    }

    public void put(K key, V value, long ttlMs) {
        store.put(key, value);
        expirations.offer(new ExpiryKey<>(key, ttlMs));
    }

    public V get(K key) {
        return store.get(key);
    }
}

실전 코드: 지수 백오프 Retry 큐

// 실패한 요청을 점점 늘어나는 간격으로 재시도
record RetryTask(String id, int attempt, Runnable job, long readyAt)
        implements Delayed {

    static RetryTask of(String id, int attempt, Runnable job) {
        long delay = (long) Math.pow(2, attempt) * 500;  // 500ms, 1s, 2s, 4s...
        return new RetryTask(id, attempt, job, System.currentTimeMillis() + delay);
    }

    @Override
    public long getDelay(TimeUnit unit) {
        return unit.convert(readyAt - System.currentTimeMillis(), MILLISECONDS);
    }

    @Override
    public int compareTo(Delayed o) {
        return Long.compare(readyAt, ((RetryTask) o).readyAt());
    }
}

DelayQueue<RetryTask> retryQueue = new DelayQueue<>();

// 소비자
Thread worker = Thread.ofVirtual().start(() -> {
    while (true) {
        RetryTask task = retryQueue.take();
        try {
            task.job().run();
        } catch (Exception e) {
            if (task.attempt() < 5) {
                retryQueue.offer(RetryTask.of(task.id(), task.attempt() + 1, task.job()));
            }
        }
    }
});

ScheduledExecutorService 와의 비교

ScheduledThreadPoolExecutor 는 내부에 DelayedWorkQueue (DelayQueue 유사 구조) 를 사용하며, 스레드 관리까지 포함.

항목DelayQueueScheduledExecutorService
스레드 관리직접자동 (thread pool)
커스텀 만료 로직✓ (Delayed 구현)제한적
반복 실행별도 구현 필요scheduleAtFixedRate 내장
취소remove(task)Future.cancel()
복잡도낮음높음

단순 delay 큐나 TTL 캐시 수준은 DelayQueue 가 더 가볍다. 주기적 작업 스케줄링은 ScheduledExecutorService 가 적합.

함정

1. Clock 비교 기준

getDelay 에서 System.currentTimeMillis()System.nanoTime() 을 혼용하면 안 된다. 일관된 시계를 사용할 것.

WARNING

System.nanoTime() 은 경과 시간 측정용, System.currentTimeMillis() 는 절대 시각용. 섞으면 overflow 나 음수 delay 가 발생할 수 있다.

2. compareTo 를 단순 빼기로 구현

// ❌ overflow 가능
public int compareTo(Delayed o) {
    return (int) (this.readyAt - o.readyAt);  // long -> int 캐스트, overflow
}

// ✓ Long.compare 사용
public int compareTo(Delayed o) {
    return Long.compare(this.readyAt, o.readyAt);
}

3. peek() 는 만료 전에도 반환

peek() 는 head 를 확인만 하고 만료 여부를 무시한다. poll() 은 만료된 경우만 꺼내고 아니면 null 반환.

DelayedTask head = q.peek();   // 만료 전이어도 반환 (null 아님)
DelayedTask ready = q.poll();  // 만료된 경우만 반환, 아니면 null
DelayedTask wait  = q.take();  // 만료될 때까지 block

CAUTION

멀티스레드 환경에서 여러 consumer 가 take() 를 동시에 기다리면 모두 정확히 만료 시각에 한 번씩 깨어난다. 하나의 원소를 여러 스레드가 가져가지 않는다. consumer 스레드 수는 처리 속도에 맞게 조정할 것.

크기 제한 없음 주의

DelayQueue 는 크기 제한이 없는 unbounded 큐. offer() 는 항상 true 를 반환하고 메모리가 허용하는 한 계속 쌓인다.

// ❌ 소비자 속도보다 생산자가 빠르면 OOM
while (running) {
    q.offer(new DelayedTask(...));  // 무한 증가 가능
}

// ✓ 외부에서 크기 모니터링
if (q.size() > MAX_PENDING) {
    log.warn("지연 큐 포화: {}", q.size());
    // 생산자 속도 조절 또는 오래된 항목 드롭
}

WARNING

DelayQueue 는 backpressure 가 없다. 대용량 스케줄링에는 ScheduledExecutorService 처럼 처리량 제어가 내장된 솔루션이 더 안전하다.

관련 위키

이 글의 용어 (6개)
[Java] ArrayBlockingQueuejava
정의 는 고정 크기 원형 배열 기반의 BlockingQueue. 생성 시 capacity 를 지정해야 하고 그 이상 늘어나지 않는다. 내부에 단일 ReentrantLock 을 사…
[Java] BlockingQueuejava
정의 는 요소를 가져올 때 비어 있으면 대기, 넣을 때 가득 차 있으면 대기 하는 thread-safe 큐 인터페이스. 생산자-소비자 (producer-consumer) 패턴의 …
[Java] ConcurrentHashMapjava
정의 는 고동시성 환경에서 사용 가능한 구현. JSR-166 (Java 5) 도입, Java 8 에서 내부 구조가 크게 재작성됐다. 과 같은 인터페이스를 제공하면서 thread-…
[Java] PriorityBlockingQueuejava
정의 는 의 동시성 + blocking 버전. unbounded 힙 기반 . - 원소는 (자연 순서) 또는 생성자에 전달한 로 정렬 - 은 절대 블록하지 않음 (unbounded…
[Java] PriorityQueuejava
정의 는 이진 힙 (binary heap) 기반의 구현. FIFO 가 아닌 우선순위 순서 로 원소를 꺼낸다. 기본은 min-heap, 의 자연 순서나 로 우선순위 결정. 가장 작…
[Java] ReentrantLockjava
정의 는 키워드와 같은 상호 배제 (mutual exclusion) 를 제공하는 클래스 기반 락. JSR-166 (Java 5) 에서 추가됐다. "재진입 (reentrant)" …

이 개념을 다룬 위키 페이지 (1)

💬 댓글

사이트 검색 / 명령어

검색

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