전통적인 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