Java 클라이언트 사용
Java 클라이언트 사용 (Java client use)
Java 클라이언트로 프로듀서, 컨슈머, 리더를 만들어 메시지를 보내고 받는 기본적인 방법을 살펴볼게요. PulsarClient 객체만 만들면 토픽을 지정해 각 개체를 만들 수 있어요. 어떤 타입의 메시지를 다루는지, 그리고 닫는 작업까지 함께 정리해볼게요.
출처: 문서
본문
프로듀서 만들기 (Create a producer)
PulsarClient 객체를 만들었다면, 특정 Pulsar 토픽에 대한 프로듀서를 만들 수 있어요.
Producer<byte[]> producer = client.newProducer()
.topic("my-topic")
.create();
// You can then send messages to the broker and topic you specified:
producer.send("My message".getBytes());
기본적으로 프로듀서는 바이트 배열로 구성된 메시지를 생산해요. 메시지 스키마를 지정하면 다른 타입을 생산할 수 있어요.
Producer<String> stringProducer = client.newProducer(Schema.STRING)
.topic("my-topic")
.create();
stringProducer.send("My message");
더 이상 필요하지 않을 때 프로듀서, 컨슈머, 클라이언트를 닫는 것을 잊지 마세요.
producer.close();
consumer.close();
client.close();
닫기 작업은 비동기로도 할 수 있어요.
producer.closeAsync()
.thenRun(() -> System.out.println("Producer closed"))
.exceptionally((ex) -> {
System.err.println("Failed to close producer: " + ex);
return null;
});
컨슈머 만들기 (Create a consumer)
Pulsar에서 컨슈머는 토픽을 구독하고, 프로듀서가 그 토픽에 게시하는 메시지를 처리해요. 먼저 PulsarClient 객체를 만들고 Pulsar 브로커의 URL을 전달해 새 컨슈머를 만들 수 있어요(위와 같이). PulsarClient 객체를 만들었다면, 토픽과 서브스크립션을 지정해 컨슈머를 만들 수 있어요.
Consumer<byte[]> consumer = client.newConsumer()
.topic("my-topic")
.subscriptionName("my-subscription")
.subscribe();
subscribe 메서드는 컨슈머를 지정된 토픽과 서브스크립션에 자동으로 구독시켜요. 컨슈머가 토픽을 듣게 하는 한 가지 방법은 while 루프를 설정하는 거예요. 이 예시 루프에서 컨슈머는 메시지를 기다리고, 받은 메시지의 내용을 출력한 뒤, 메시지가 처리되었음을 ack해요. 처리 로직이 실패하면 부정 확인(negative acknowledgment)을 사용해 메시지를 나중에 재전달받을 수 있어요.
while (true) {
// Wait for a message
Message<byte[]> msg = consumer.receive();
try {
// Do something with the message
System.out.println("Message received: " + new String(msg.getData()));
// Acknowledge the message
consumer.acknowledge(msg);
} catch (Exception e) {
// Message failed to process, redeliver later
consumer.negativeAcknowledge(msg);
}
}
메인 스레드를 블록하지 않고 새 메시지를 계속 듣고 싶다면 MessageListener를 사용하는 것을 고려해보세요. MessageListener는 PulsarClient 안의 스레드 풀을 사용해요. 메시지 리스너에 사용할 스레드 수는 ClientBuilder에서 설정할 수 있어요.
MessageListener<byte[]> myMessageListener = (consumer, msg) -> {
try {
System.out.println("Message received: " + new String(msg.getData()));
consumer.acknowledge(msg);
} catch (Exception e) {
consumer.negativeAcknowledge(msg);
}
};
Consumer<byte[]> consumer = client.newConsumer()
.topic("my-topic")
.subscriptionName("my-subscription")
.messageListener(myMessageListener)
.subscribe();
리더 만들기 (Create a reader)
리더 인터페이스를 사용하면 Pulsar 클라이언트가 토픽 안에서 "수동으로 위치를 잡고"(manually position) 지정된 메시지부터 이후의 모든 메시지를 읽을 수 있어요. Java용 Pulsar API는 토픽과 MessageId를 지정해 Reader 객체를 만들 수 있게 해줘요. 다음은 예시예요.
byte[] msgIdBytes = // Some message ID byte array
MessageId id = MessageId.fromByteArray(msgIdBytes);
Reader<byte[]> reader = pulsarClient.newReader()
.topic(topic)
.startMessageId(id)
.create();
while (true) {
Message<byte[]> message = reader.readNext();
// Process message
}
위 예시에서 Reader 객체는 특정 토픽과 메시지(ID로)에 대해 인스턴스화돼요. 리더는 msgIdBytes가 식별하는 메시지 이후의 각 메시지를 반복해 읽어요(그 값을 어떻게 얻는지는 애플리케이션에 달려 있어요).
위 코드는 Reader 객체를 특정 메시지(ID로)를 가리키게 하는 모습을 보여줘요. 하지만 MessageId.earliest를 사용해 토픽에서 가장 이른 메시지를 가리키거나, MessageId.latest로 가장 최근의 사용 가능한 메시지를 가리킬 수도 있어요.
더 알아보기 (Learn more)
- Java 클라이언트 라이브러리 설정 — Java 클라이언트를 설치하는 방법을 알아봐요.
- Java 클라이언트 초기화 — PulsarClient를 만드는 방법을 살펴봐요.
- Java 클라이언트 — Java 클라이언트 개요를 확인해요.
- 프로듀서 사용하기 — 프로듀서 개념을 자세히 알아봐요.
- 컨슈머 사용하기 — 컨슈머 개념을 자세히 알아봐요.
- Java 클라이언트(V5) — 스케일러블 토픽용 V5 SDK를 살펴봐요.