블로킹 큐를 활용한 생산자-소비자 패턴 구현

전통적인 wait-notify 메커니즘을 활용한 동기화는 코드를 복잡하게 만들기 쉽습니다. java.util.concurrent 패키지의 BlockingQueue 인터페이스는 이러한 동기화 문제를 간결하게 해결해 줍니다. 생산자-소비자(Producer-Consumer) 패턴의 핵심은 자원을 생성하는 생산자의 put() 호출과 자원을 소비하는 소비자의 take() 호출을 제어하는 것입니다. BlockingQueue는 큐가 가득 찼을 때 put()을 호출한 스레드를 대기시키고, 큐가 비어있을 때 take()을 호출한 스레드를 대기시키는 블로킹 제어를 기본적으로 제공하므로, 별도의 동기화 로직 없이도 스레드 안전성을 보장합니다.

생산자 스레드 구현

아래는 데이터를 생성하여 큐에 삽입하는 생산자 클래스입니다.

import java.util.concurrent.BlockingQueue;

public class DataProducer implements Runnable {
    private final BlockingQueue<String> sharedBuffer;

    public DataProducer(BlockingQueue<String> buffer) {
        this.sharedBuffer = buffer;
    }

    @Override
    public void run() {
        int itemId = 0;
        try {
            while (!Thread.currentThread().isInterrupted()) {
                String producedItem = generateItem(itemId++);
                sharedBuffer.put(producedItem);
                System.out.println("생산 완료: " + producedItem + " | 현재 버퍼 크기: " + sharedBuffer.size());
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            System.out.println("생산자 스레드가 중단되었습니다.");
        }
    }

    private String generateItem(int id) {
        try {
            Thread.sleep(50); // 자원 생성 지연 시뮬레이션
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
        return "Item-" + id;
    }
}

생산자 스레드는 자원을 생성하여 큐에 넣습니다. 만약 큐의 용량이 최대치(예: 15)에 도달했다면, 새로운 공간이 확보될 때까지 put() 메서드에서 블로킹되므로 큐의 크기는 설정된 최댓값을 초과하지 않습니다.

소비자 스레드 구현

아래는 큐에서 데이터를 꺼내어 소비하는 소비자 클래스입니다.

import java.util.concurrent.BlockingQueue;

public class DataConsumer implements Runnable {
    private final BlockingQueue<String> sharedBuffer;

    public DataConsumer(BlockingQueue<String> buffer) {
        this.sharedBuffer = buffer;
    }

    @Override
    public void run() {
        try {
            while (!Thread.currentThread().isInterrupted()) {
                String consumedItem = sharedBuffer.take();
                System.out.println("소비 완료: " + consumedItem + " | 현재 버퍼 크기: " + sharedBuffer.size());
                processItem(consumedItem);
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            System.out.println("소비자 스레드가 중단되었습니다.");
        }
    }

    private void processItem(String item) {
        try {
            Thread.sleep(100); // 자원 처리 지연 시뮬레이션
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

소비자 스레드는 큐에 데이터가 존재할 때 take()를 통해 값을 꺼내 처리합니다. 큐가 비어있다면 생산자가 데이터를 삽입할 때까지 대기 상태에 진입합니다.

동작 확인

구현된 생산자와 소비자 스레드를 실행하여 패턴이 정상적으로 동작하는지 테스트합니다.

import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;

public class BlockingQueueDemo {
    public static void main(String[] args) throws InterruptedException {
        int producerCount = 3;
        int consumerCount = 2;

        BlockingQueue<String> blockingBuffer = new LinkedBlockingQueue<>(15);

        for (int i = 0; i < producerCount; i++) {
            new Thread(new DataProducer(blockingBuffer), "Producer-" + i).start();
        }

        for (int i = 0; i < consumerCount; i++) {
            new Thread(new DataConsumer(blockingBuffer), "Consumer-" + i).start();
        }

        Thread.sleep(5000);
        System.exit(0);
    }
}

위 코드를 실행하면 다음과 유사한 로그를 확인할 수 있으며, 버퍼의 크기가 설정된 최댓값(15)을 넘지 않고 생산과 소비가 안정적으로 이루어지는 것을 볼 수 있습니다.

생산 완료: Item-0 | 현재 버퍼 크기: 1
소비 완료: Item-0 | 현재 버퍼 크기: 0
생산 완료: Item-1 | 현재 버퍼 크기: 1
생산 완료: Item-2 | 현재 버퍼 크기: 2
소비 완료: Item-1 | 현재 버퍼 크기: 1
...
생산 완료: Item-15 | 현재 버퍼 크기: 15
소비 완료: Item-12 | 현재 버퍼 크기: 14
생산 완료: Item-16 | 현재 버퍼 크기: 15
소비 완료: Item-13 | 현재 버퍼 크기: 14

태그: java BlockingQueue Producer-Consumer concurrency LinkedBlockingQueue

8월 15일 17:06에 게시됨