리더 사용하기
리더 사용하기 (Readers)
리더는 커서 없는 컨슈머예요. 즉 Pulsar가 진행 상황을 추적하지 않고 메시지를 ack할 필요가 없어요. 이 문서는 토픽의 가장 이른 메시지부터 읽는 법, 특정 메시지부터 읽는 법, 청킹과 인터셉터 구성, 스티키 키 레인지 리더까지 정리했어요.
출처: 문서
본문
클라이언트를 설정한 뒤에는 더 살펴보며 리더로 작업을 시작할 수 있어요.
메시지 수신 및 읽기 (Receive and read messages)
리더는 커서 없는 컨슈머일 뿐이에요. 즉 Pulsar가 진행 상황을 추적하지 않고 메시지를 ack할 필요가 없어요.
다음은 토픽에서 가장 이른 사용 가능 메시지부터 읽기 시작하는 예시예요.
- Java C#
import org.apache.pulsar.client.api.Message;
import org.apache.pulsar.client.api.MessageId;
import org.apache.pulsar.client.api.Reader;
// Create a reader on a topic and for a specific message (and onward)
Reader<byte[]> reader = pulsarClient.newReader()
.topic("reader-api-test")
.startMessageId(MessageId.earliest)
.create();
while (true) {
Message message = reader.readNext();
// Process the message
}
await foreach (var message in reader.Messages())
{
Console.WriteLine("Received: " + Encoding.UTF8.GetString(message.Data.ToArray()));
}
다음 메시지 읽기 (Read next message)
가장 최근 사용 가능 메시지부터 읽는 리더를 만들려면:
- Java Go
Reader<byte[]> reader = pulsarClient.newReader()
.topic(topic)
.startMessageId(MessageId.latest)
.create();
import (
"context"
"fmt"
"log"
"github.com/apache/pulsar-client-go/pulsar"
)
func main() {
client, err := pulsar.NewClient(pulsar.ClientOptions{URL: "pulsar://localhost:6650"})
if err != nil {
log.Fatal(err)
}
defer client.Close()
reader, err := client.CreateReader(pulsar.ReaderOptions{
Topic: "topic-1",
StartMessageID: pulsar.EarliestMessageID(),
})
if err != nil {
log.Fatal(err)
}
defer reader.Close()
for reader.HasNext() {
msg, err := reader.Next(context.Background())
if err != nil {
log.Fatal(err)
}
fmt.Printf("Received message msgId: %#v -- content: '%s'\n",
msg.ID(), string(msg.Payload()))
}
}
위 예시에서 리더는 가장 이른 사용 가능 메시지(pulsar.EarliestMessage로 지정)부터 읽기 시작해요. 리더는 pulsar.LatestMessage의 가장 최근 메시지나, DeserializeMessageID 함수(bytes 배열을 받아 MessageID 객체를 반환)로 지정한 다른 메시지 ID부터 읽을 수도 있어요. 예시:
lastSavedId := // Read last saved message id from external store as byte[]
reader, err := client.CreateReader(pulsar.ReaderOptions{
Topic: "my-golang-topic",
StartMessageID: pulsar.DeserializeMessageID(lastSavedId),
})
특정 메시지 읽기 (Read specific messages)
가장 이른 메시지와 가장 최근 메시지 사이의 어떤 메시지부터 읽는 리더를 만들려면:
- Java Go
byte[] msgIdBytes = // Some byte array
MessageId id = MessageId.fromByteArray(msgIdBytes);
Reader<byte[]> reader = pulsarClient.newReader()
.topic(topic)
.startMessageId(id)
.create();
client, err := pulsar.NewClient(pulsar.ClientOptions{
URL: "pulsar://localhost:6650",
})
if err != nil {
log.Fatal(err)
}
defer client.Close()
topic := "topic-1"
ctx := context.Background()
// create producer
producer, err := client.CreateProducer(pulsar.ProducerOptions{
Topic: topic,
DisableBatching: true,
})
if err != nil {
log.Fatal(err)
}
defer producer.Close()
// send 10 messages
msgIDs := [10]pulsar.MessageID{}
for i := 0; i < 10; i++ {
msgID, _ := producer.Send(ctx, &pulsar.ProducerMessage{
Payload: []byte(fmt.Sprintf("hello-%d", i)),
})
msgIDs[i] = msgID
}
// create reader on 5th message (not included)
reader, err := client.CreateReader(pulsar.ReaderOptions{
Topic: topic,
StartMessageID: msgIDs[4],
StartMessageIDInclusive: false,
})
if err != nil {
log.Fatal(err)
}
defer reader.Close()
// receive the remaining 5 messages
for i := 5; i < 10; i++ {
msg, err := reader.Next(context.Background())
if err != nil {
log.Fatal(err)
}
fmt.Printf("Read %d-th msg: %s\n", i, string(msg.Payload()))
}
// create reader on 5th message (included)
readerInclusive, err := client.CreateReader(pulsar.ReaderOptions{
Topic: topic,
StartMessageID: msgIDs[4],
StartMessageIDInclusive: true,
})
if err != nil {
log.Fatal(err)
}
defer readerInclusive.Close()
청킹 구성 (Configure chunking)
리더의 청킹 구성은 컨슈머와 비슷해요. 자세한 내용은 컨슈머 청킹 구성을 참고하세요.
다음은 리더에 메시지 청킹을 구성하는 예시예요.
- Java
Reader<byte[]> reader = pulsarClient.newReader()
.topic(topicName)
.startMessageId(MessageId.earliest)
.maxPendingChunkedMessage(12)
.autoAckOldestChunkedMessageOnQueueFull(true)
.expireTimeOfIncompleteChunkedMessage(12, TimeUnit.MILLISECONDS)
.create();
메시지 인터셉트 (Intercept messages)
Pulsar 리더 인터셉터는 Pulsar 리더가 메시지를 읽기 전에 사용자 정의 처리로 메시지를 가로채고 변형할 수 있어요. 리더 인터셉터를 사용하면 메시지 수정, 속성 추가, 통계 수집 등을 각각 유사한 메커니즘을 만들지 않고도 메시지를 읽기 전에 통일된 메시징 처리를 적용할 수 있어요.
Pulsar 리더 인터셉터는 Pulsar 컨슈머 인터셉터 위에서 동작해요. 플러그인 인터페이스 ReaderInterceptor는 ConsumerInterceptor의 하위 집합으로 볼 수 있고 두 가지 주요 이벤트가 있어요.
beforeRead— 리더가 메시지를 읽기 전에 트리거돼요. 이 이벤트에서 메시지를 수정할 수 있어요.onPartitionsChange— 파티션에 변경이 감지될 때 트리거돼요.
트리거된 이벤트를 감지하고 커스텀 처리를 수행하려면 Reader를 만들 때 ReaderInterceptor를 추가할 수 있어요.
- Java
PulsarClient pulsarClient = PulsarClient.builder().serviceUrl("pulsar://localhost:6650").build();
Reader<byte[]> reader = pulsarClient.newReader()
.topic("t1")
.autoUpdatePartitionsInterval(5, TimeUnit.SECONDS)
.intercept(new ReaderInterceptor<byte[]>() {
@Override
public void close() {
}
@Override
public Message<byte[]> beforeRead(Reader<byte[]> reader, Message<byte[]> message) {
// user-defined processing logic
return message;
}
@Override
public void onPartitionsChange(String topicName, int partitions) {
// user-defined processing logic
}
})
.startMessageId(MessageId.earliest)
.create();
스티키 키 레인지 리더 (Sticky key range reader)
스티키 키 레인지 리더에서 브로커는 메시지 키의 해시가 지정된 키 해시 레인지에 포함된 메시지만 전달해요. 한 리더에 여러 키 해시 레인지를 지정할 수 있어요.
다음은 스티키 키 레인지 리더를 만드는 예시예요.
- Java
pulsarClient.newReader()
.topic(topic)
.startMessageId(MessageId.earliest)
.keyHashRange(Range.of(0, 10000), Range.of(20001, 30000))
.create();
전체 해시 레인지 크기는 65536이므로, 레인지의 최대 끝값은 65535보다 작거나 같아야 해요.
더 알아보기 (Learn more)
- 프로듀서 사용하기 — 메시지를 게시하는 방법을 알아봐요.
- 컨슈머 사용하기 — 메시지를 소비하는 방법을 알아봐요.
- 클라이언트 라이브러리 — 다양한 언어 클라이언트를 살펴봐요.
- 리더 개념 — 리더가 무엇인지 자세히 알아봐요.