리더 사용하기

리더 사용하기 (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()));
}

가장 최근 사용 가능 메시지부터 읽는 리더를 만들려면:

  • 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 컨슈머 인터셉터 위에서 동작해요. 플러그인 인터페이스 ReaderInterceptorConsumerInterceptor의 하위 집합으로 볼 수 있고 두 가지 주요 이벤트가 있어요.

  • 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)