컨슈머 사용하기

컨슈머 사용하기 (Consumers)

클라이언트를 설정했다면, 이제 컨슈머로 메시지를 받는 작업을 시작할 수 있어요. Pulsar 컨슈머는 토픽 구독, 메시지 수신 방식, ack 방식이 다양해서 상황에 맞게 골라 쓸 수 있죠. 이 문서에서 서브스크립션 타입별 차이부터 메시지 수신, ack, 청킹, 리스너, 인터셉터까지 하나씩 살펴볼게요.

출처: 문서

본문

클라이언트를 설정한 뒤에는 컨슈머 작업을 시작하기 위해 더 살펴볼 수 있어요.

토픽 구독 (Subscribe to topics)

Pulsar는 다양한 시나리오에 맞는 여러 서브스크립션 타입을 제공해요. 하나의 토픽은 서로 다른 서브스크립션 타입의 여러 서브스크립션을 가질 수 있어요. 하지만 하나의 서브스크립션은 한 번에 하나의 서브스크립션 타입만 가질 수 있어요.

서브스크립션은 서브스크립션 이름과 동일해요. 서브스크립션 이름은 한 번에 하나의 서브스크립션 타입만 지정할 수 있어요. 서브스크립션 타입을 바꾸려면 먼저 이 서브스크립션의 모든 컨슈머를 중지해야 해요.

서브스크립션 타입마다 메시지 배포(distribution) 방식이 달라요. 이 섹션에서는 서브스크립션 타입 간 차이와 사용 방법을 설명해요.

차이를 잘 설명하기 위해 "my-topic"이라는 토픽이 있고, 프로듀서가 10개의 메시지를 게시했다고 가정할게요.

Producer<String> producer = client.newProducer(Schema.STRING)
        .topic("my-topic")
        .enableBatching(false)
        .create();
// 3 messages with "key-1", 3 messages with "key-2", 2 messages with "key-3" and 2 messages with "key-4"
producer.newMessage().key("key-1").value("message-1-1").send();
producer.newMessage().key("key-1").value("message-1-2").send();
producer.newMessage().key("key-1").value("message-1-3").send();
producer.newMessage().key("key-2").value("message-2-1").send();
producer.newMessage().key("key-2").value("message-2-2").send();
producer.newMessage().key("key-2").value("message-2-3").send();
producer.newMessage().key("key-3").value("message-3-1").send();
producer.newMessage().key("key-3").value("message-3-2").send();
producer.newMessage().key("key-4").value("message-4-1").send();
producer.newMessage().key("key-4").value("message-4-2").send();
producer = client.create_producer('my-topic', batching_enabled=False)
# 3 messages with "key-1", 3 messages with "key-2", 2 messages with "key-3" and 2 messages with "key-4"
producer.send(b'message-1-1', partition_key='key-1')
producer.send(b'message-1-2', partition_key='key-1')
producer.send(b'message-1-3', partition_key='key-1')
producer.send(b'message-2-1', partition_key='key-2')
producer.send(b'message-2-2', partition_key='key-2')
producer.send(b'message-2-3', partition_key='key-2')
producer.send(b'message-3-1', partition_key='key-3')
producer.send(b'message-3-2', partition_key='key-3')
producer.send(b'message-4-1', partition_key='key-4')
producer.send(b'message-4-2', partition_key='key-4')

Exclusive (독점)

새 컨슈머를 만들고 Exclusive 서브스크립션 타입으로 구독해요.

Consumer<byte[]> consumer = client.newConsumer()
        .topic("my-topic")
        .subscriptionName("my-subscription")
        .subscriptionType(SubscriptionType.Exclusive)
        .subscribe();
consumer = client.subscribe('my-topic', 'my-subscription',
                             consumer_type=pulsar.ConsumerType.Exclusive)

첫 번째 컨슈머만 서브스크립션에 허용되고, 다른 컨슈머는 에러를 받아요. 첫 번째 컨슈머는 10개의 메시지를 모두 받으며, 소비 순서는 생산 순서와 같아요. 토픽이 파티셔닝되어 있다면, 첫 번째 컨슈머가 모든 파티션 토픽을 구독하고 다른 컨슈머는 파티션을 할당받지 못해 에러를 받아요.

Failover (장애 조치)

새 컨슈머를 만들고 Failover 서브스크립션 타입으로 구독해요.

Consumer<byte[]> consumer1 = client.newConsumer()
        .topic("my-topic")
        .subscriptionName("my-subscription")
        .subscriptionType(SubscriptionType.Failover)
        .subscribe();
Consumer<byte[]> consumer2 = client.newConsumer()
        .topic("my-topic")
        .subscriptionName("my-subscription")
        .subscriptionType(SubscriptionType.Failover)
        .subscribe();
//consumer1 is the active consumer, consumer2 is the standby consumer.
//consumer1 receives 5 messages and then crashes, consumer2 takes over as an active consumer.
consumer1 = client.subscribe('my-topic', 'my-subscription',
                              consumer_type=pulsar.ConsumerType.Failover)
consumer2 = client.subscribe('my-topic', 'my-subscription',
                              consumer_type=pulsar.ConsumerType.Failover)

여러 컨슈머가 같은 서브스크립션에 연결될 수 있지만, 첫 번째 컨슈머만 활성(active)이고 나머지는 대기(standby) 상태예요. 활성 컨슈머가 연결이 끊기면, 메시지가 대기 컨슈머 중 하나로 전달되고 그 컨슈머가 활성 컨슈머가 돼요.

첫 번째 활성 컨슈머가 5개의 메시지를 받은 후 연결이 끊기면, 대기 컨슈머가 활성 컨슈머가 돼요. Consumer1은 다음을 받게 돼요:

("key-1", "message-1-1")
("key-1", "message-1-2")
("key-1", "message-1-3")
("key-2", "message-2-1")
("key-2", "message-2-2")

consumer2는 다음을 받게 돼요:

("key-2", "message-2-3")
("key-3", "message-3-1")
("key-3", "message-3-2")
("key-4", "message-4-1")
("key-4", "message-4-2")

토픽이 파티셔닝된 토픽이라면, 각 파티션에는 활성 컨슈머가 하나만 있고, 한 파티션의 메시지는 한 컨슈머에게만 배포되며, 여러 파티션의 메시지는 여러 컨슈머에게 배포돼요.

Shared (공유)

새 컨슈머를 만들고 Shared 서브스크립션 타입으로 구독해요.

Consumer<byte[]> consumer1 = client.newConsumer()
        .topic("my-topic")
        .subscriptionName("my-subscription")
        .subscriptionType(SubscriptionType.Shared)
        .subscribe();

Consumer<byte[]> consumer2 = client.newConsumer()
        .topic("my-topic")
        .subscriptionName("my-subscription")
        .subscriptionType(SubscriptionType.Shared)
        .subscribe();
//Both consumer1 and consumer2 are active consumers.
consumer1 = client.subscribe('my-topic', 'my-subscription',
                              consumer_type=pulsar.ConsumerType.Shared)
consumer2 = client.subscribe('my-topic', 'my-subscription',
                              consumer_type=pulsar.ConsumerType.Shared)

Shared 서브스크립션 타입에서는 여러 컨슈머가 같은 서브스크립션에 연결될 수 있고, 메시지는 컨슈머들 사이에 라운드 로빈(round-robin) 방식으로 배포돼요. 브로커가 한 번에 메시지 하나만 배포한다면, consumer1은 다음을 받아요:

("key-1", "message-1-1")
("key-1", "message-1-3")
("key-2", "message-2-2")
("key-3", "message-3-1")
("key-4", "message-4-1")

consumer2는 다음을 받아요:

("key-1", "message-1-2")
("key-2", "message-2-1")
("key-2", "message-2-3")
("key-3", "message-3-2")
("key-4", "message-4-2")

Shared 서브스크립션은 Exclusive 및 Failover 타입과 달라요. Shared는 유연성이 좋지만 순서 보장을 제공하지 못해요.

Key_Shared (키 공유)

이것은 2.4.0 릴리스부터 도입된 새 서브스크립션 타입이에요. 새 컨슈머를 만들고 Key_Shared 서브스크립션 타입으로 구독해요.

Key_Shared 서브스크립션을 사용할 때는 프로듀서가 배칭을 비활성화하거나 키 기반 배칭(예: Java의 BatcherBuilder.KEY_BASED)을 사용해야 해요. 기본 배칭은 서로 다른 키를 가진 메시지를 같은 배치에 넣어 Key_Shared 라우팅 의미를 깨뜨릴 수 있어요. 아래 코드 예시를 참고하세요.

Consumer<byte[]> consumer1 = client.newConsumer()
        .topic("my-topic")
        .subscriptionName("my-subscription")
        .subscriptionType(SubscriptionType.Key_Shared)
        .subscribe();

Consumer<byte[]> consumer2 = client.newConsumer()
        .topic("my-topic")
        .subscriptionName("my-subscription")
        .subscriptionType(SubscriptionType.Key_Shared)
        .subscribe();
//Both consumer1 and consumer2 are active consumers.
consumer1 = client.subscribe('my-topic', 'my-subscription',
                              consumer_type=pulsar.ConsumerType.KeyShared)
consumer2 = client.subscribe('my-topic', 'my-subscription',
                              consumer_type=pulsar.ConsumerType.KeyShared)

Shared 서브스크립션과 마찬가지로, Key_Shared 타입의 모든 컨슈머는 같은 서브스크립션에 연결될 수 있어요. 하지만 Key_Shared 타입은 Shared와 달라요. Key_Shared 타입에서는 같은 키를 가진 메시지가 순서대로 한 컨슈머에게만 전달돼요. 서로 다른 컨슈머 간의 가능한 메시지 배포(기본적으로 어떤 키가 어떤 컨슈머에 할당될지 미리 알 수 없지만, 한 키는 동시에 한 컨슈머에게만 할당돼요):

consumer1은 다음을 받아요:

("key-1", "message-1-1")
("key-1", "message-1-2")
("key-1", "message-1-3")
("key-3", "message-3-1")
("key-3", "message-3-2")

consumer2는 다음을 받아요:

("key-2", "message-2-1")
("key-2", "message-2-2")
("key-2", "message-2-3")
("key-4", "message-4-1")
("key-4", "message-4-2")

프로듀서 쪽에서 배칭이 활성화되면, 기본적으로 서로 다른 키를 가진 메시지가 한 배치에 추가돼요. 브로커는 그 배치를 컨슈머에게 배포하므로, 기본 배치 메커니즘이 Key_Shared 서브스크립션에서 보장하는 메시지 배포 의미를 깨뜨릴 수 있어요. 프로듀서는 KeyBasedBatcher를 사용해야 해요.

Producer producer = client.newProducer()
        .topic("my-topic")
        .batcherBuilder(BatcherBuilder.KEY_BASED)
        .create();
producer = client.create_producer('my-topic',
                                   batching_type=pulsar.BatchingType.KeyBased)

또는 프로듀서가 배칭을 비활성화할 수 있어요.

Producer producer = client.newProducer()
        .topic("my-topic")
        .enableBatching(false)
        .create();
producer = client.create_producer('my-topic', batching_enabled=False)

메시지 키가 지정되지 않으면, 기본적으로 키가 없는 메시지는 순서대로 한 컨슈머에게 배포돼요.

여러 토픽 구독 (Subscribe to multi-topics)

컨슈머를 단일 Pulsar 토픽에 구독하는 것 외에도, 다중 토픽 서브스크립션(multi-topic subscriptions)을 사용해 여러 토픽을 동시에 구독할 수 있어요. 다중 토픽 서브스크립션을 사용하려면 정규식(regex) 또는 토픽 List를 제공할 수 있어요. regex로 토픽을 선택하면 모든 토픽이 같은 Pulsar 네임스페이스 안에 있어야 해요.

다음은 몇 가지 예시예요.

import org.apache.pulsar.client.api.Consumer;
import org.apache.pulsar.client.api.PulsarClient;

import java.util.Arrays;
import java.util.List;
import java.util.regex.Pattern;

ConsumerBuilder consumerBuilder = pulsarClient.newConsumer()
        .subscriptionName(subscription);

// Subscribe to all topics in a namespace
Pattern allTopicsInNamespace = Pattern.compile("public/default/.*");
Consumer allTopicsConsumer = consumerBuilder
        .topicsPattern(allTopicsInNamespace)
        .subscribe();

// Subscribe to a subsets of topics in a namespace, based on regex
Pattern someTopicsInNamespace = Pattern.compile("public/default/foo.*");
Consumer allTopicsConsumer = consumerBuilder
        .topicsPattern(someTopicsInNamespace)
        .subscribe();

위 예시에서 컨슈머는 토픽 이름 패턴에 일치하는 영속 토픽을 구독해요. 컨슈머가 토픽 이름 패턴에 일치하는 모든 영속 및 비영속 토픽을 구독하려면 subscriptionTopicsModeRegexSubscriptionMode.AllTopics로 설정하세요.

Pattern pattern = Pattern.compile("public/default/.*");
pulsarClient.newConsumer()
    .subscriptionName("my-sub")
    .topicsPattern(pattern)
    .subscriptionTopicsMode(RegexSubscriptionMode.AllTopics)
    .subscribe();

참고: 기본적으로 컨슈머의 subscriptionTopicsModePersistentOnly예요. subscriptionTopicsMode의 사용 가능한 옵션은 PersistentOnly, NonPersistentOnly, AllTopics예요.

명시적인 토픽 목록(원한다면 네임스페이스 건너편)을 구독할 수도 있어요:

List<String> topics = Arrays.asList("topic-1", "topic-2", "topic-3");
Consumer multiTopicConsumer = consumerBuilder
        .topics(topics)
        .subscribe();

// Alternatively:
Consumer multiTopicConsumer = consumerBuilder
        .topic("topic-1", "topic-2", "topic-3")
        .subscribe();

동기 subscribe 메서드 대신 subscribeAsync 메서드로 여러 토픽을 비동기로 구독할 수도 있어요. 다음은 예시예요.

Pattern allTopicsInNamespace = Pattern.compile("persistent://public/default.*");
consumerBuilder.topics(topics).subscribeAsync().thenAccept(this::receiveMessageFromConsumer);

private void receiveMessageFromConsumer(Object consumer) {
    ((Consumer) consumer).receiveAsync().thenAccept(message -> {
        // Do something with the received message
        receiveMessageFromConsumer(consumer);
    });
}
client, err := pulsar.NewClient(pulsar.ClientOptions{
    URL: "pulsar://localhost:6650",
})
if err != nil {
    log.Fatal(err)
}

topics := []string{"topic-1", "topic-2"}
consumer, err := client.Subscribe(pulsar.ConsumerOptions{
    // fill `Topics` field will create a multi-topic consumer
    Topics:           topics,
    SubscriptionName: "multi-topic-sub",
})
if err != nil {
    log.Fatal(err)
}
defer consumer.Close()
import re
consumer = client.subscribe(re.compile('persistent://public/default/topic-*'), 'my-subscription')
while True:
    msg = consumer.receive()
    try:
        print("Received message '{}' id='{}'".format(msg.data(), msg.message_id()))
        # Acknowledge successful processing of the message
        consumer.acknowledge(msg)
    except Exception:
        # Message failed to be processed
        consumer.negative_acknowledge(msg)
client.close()

토픽 구독 해제 (Unsubscribe from topics)

이 예시는 컨슈머가 토픽에서 구독을 해제하는 방법을 보여줘요.

consumer.unsubscribe();
await consumer.Unsubscribe();
consumer.unsubscribe()

컨슈머가 토픽에서 구독을 해제하면 더 이상 사용할 수 없고 폐기(dispose)돼요.

메시지 수신 (Receive messages)

이 예시는 컨슈머가 토픽에서 메시지를 받는 방법을 보여줘요.

Message message = consumer.receive();
await foreach (var message in consumer.Messages())
{
    Console.WriteLine("Received: " + Encoding.UTF8.GetString(message.Data.ToArray()));
}
msg = consumer.receive()

타임아웃을 포함한 메시지 수신 (Receive messages with timeout)

consumer.receive(10, TimeUnit.SECONDS);
client, err := pulsar.NewClient(pulsar.ClientOptions{
    URL: "pulsar://localhost:6650",
})
if err != nil {
    log.Fatal(err)
}
defer client.Close()

topic := "test-topic-with-no-messages"
ctx, cancel := context.WithTimeout(context.Background(), 500*time.Millisecond)
defer cancel()

// create consumer
consumer, err := client.Subscribe(pulsar.ConsumerOptions{
    Topic:            topic,
    SubscriptionName: "my-sub1",
    Type:             pulsar.Shared,
})
if err != nil {
    log.Fatal(err)
}
defer consumer.Close()

// receive message with a timeout
msg, err := consumer.Receive(ctx)
if err != nil {
    log.Fatal(err)
}
fmt.Println(msg.Payload())
# Receive with 10 second timeout (timeout in milliseconds)
try:
    msg = consumer.receive(timeout_millis=10000)
except Exception:
    print("No message received within timeout period")

비동기 메시지 수신 (Async receive messages)

receive 메서드는 메시지를 동기적으로 받아요(메시지를 사용할 수 있을 때까지 컨슈머 프로세스가 블록돼요). 또한 비동기 receive를 사용할 수도 있는데, 새 메시지가 사용 가능해지면 즉시 CompletableFuture 객체를 반환해요. 다음은 예시예요.

CompletableFuture<Message> asyncMessage = consumer.receiveAsync();

비동기 수신 작업은 CompletableFuture 안에 감싸진 Message를 반환해요.

import asyncio

async def receive_messages():
    msg = await consumer.receive_async()
    return msg

# Use in async context
msg = asyncio.run(receive_messages())

Python의 비동기 수신 작업은 asyncio를 사용하며 Message로 해석되는 코루틴을 반환해요.

배치 메시지 수신 (Batch receive messages)

batchReceive를 사용해 호출할 때마다 여러 메시지를 받을 수 있어요. 다음은 예시예요.

Messages messages = consumer.batchReceive();
for (Object message : messages) {
  // do something
}
consumer.acknowledge(messages)
messages = consumer.batch_receive()
for msg in messages:
    # do something
    pass
consumer.acknowledge(messages)

배치 수신 정책(batch receive policy)은 단일 배치의 메시지 수와 바이트를 제한해요. 충분한 메시지를 기다릴 타임아웃을 지정할 수 있어요. 다음 조건 중 하나라도 충족되면 배치 수신이 완료돼요: 충분한 메시지 수, 메시지 바이트, 대기 타임아웃.

Consumer consumer = client.newConsumer()
    .topic("my-topic")
    .subscriptionName("my-subscription")
    .batchReceivePolicy(BatchReceivePolicy.builder()
        .maxNumMessages(100)
        .maxNumBytes(1024 * 1024)
        .timeout(200, TimeUnit.MILLISECONDS)
        .build())
    .subscribe();

기본 배치 수신 정책은 다음과 같아요:

BatchReceivePolicy.builder()
    .maxNumMessage(-1)
    .maxNumBytes(10 * 1024 * 1024)
    .timeout(100, TimeUnit.MILLISECONDS)
    .build();

메시지 ack (Acknowledge messages)

메시지는 개별적으로 또는 누적적으로 ack할 수 있어요. 메시지 ack에 대한 자세한 내용은 acknowledgment를 참고하세요.

개별 ack (Acknowledge messages individually)

consumer.acknowledge(msg);
await consumer.Acknowledge(message);
consumer.acknowledge(msg)

누적 ack (Acknowledge messages cumulatively)

consumer.acknowledgeCumulative(msg);
await consumer.AcknowledgeCumulative(message);
consumer.acknowledge_cumulative(msg)

부정 ack 재전송 백오프 (Negative acknowledgment redelivery backoff)

RedeliveryBackoff는 재전송 백오프 메커니즘을 제공해요. 메시지의 재전송 횟수를 설정해 재전송 지연을 다르게 할 수 있어요.

Consumer consumer =  client.newConsumer()
        .topic("my-topic")
        .subscriptionName("my-subscription")
        .negativeAckRedeliveryBackoff(MultiplierRedeliveryBackoff.builder()
                .minDelayMs(1000)
                .maxDelayMs(60 * 1000)
                .build())
        .subscribe();
consumer = client.subscribe(
    'my-topic',
    'my-subscription',
    negative_ack_redelivery_delay_ms=1000
)

ack 타임아웃 재전송 백오프 (Acknowledgment timeout redelivery backoff)

RedeliveryBackoff는 재전송 백오프 메커니즘을 제공해요. 메시지를 재시도하는 횟수를 설정해 재전송 지연을 다르게 할 수 있어요.

Consumer consumer =  client.newConsumer()
        .topic("my-topic")
        .subscriptionName("my-subscription")
        .ackTimeout(10, TimeUnit.SECOND)
        .ackTimeoutRedeliveryBackoff(MultiplierRedeliveryBackoff.builder()
                .minDelayMs(1000)
                .maxDelayMs(60000)
                .multiplier(2)
                .build())
        .subscribe();
consumer = client.subscribe(
    'my-topic',
    'my-subscription',
    unacked_messages_timeout_ms=10000
)

메시지 재전송 동작은 다음과 같아야 해요.

재전송 횟수 재전송 지연
1 10 + 1 seconds
2 10 + 2 seconds
3 10 + 4 seconds
4 10 + 8 seconds
5 10 + 16 seconds
6 10 + 32 seconds
7 10 + 60 seconds
8 10 + 60 seconds
  • negativeAckRedeliveryBackoffconsumer.negativeAcknowledge(MessageId messageId)와 함께 동작하지 않아요. 메시지 ID에서 재전송 횟수를 얻을 수 없기 때문이에요.
  • 컨슈머가 크래시하면, ack되지 않은 메시지의 재전송이 트리거돼요. 이 경우 RedeliveryBackoff가 적용되지 않고, 백오프의 지연 시간보다 일찍 메시지가 재전송될 수 있어요.

청킹 구성 (Configure chunking)

컨슈머가 동시에 유지하는 청크(chunk) 메시지의 최대 수를 특정 파라미터를 구성해 제한할 수 있어요. 구성된 임계값에 도달하면, 컨슈머는 보류 중인 메시지를 조용히 ack하거나 나중에 재전송하도록 브로커에 요청해 버려요. 다음은 메시지 청킹을 구성하는 예시예요.

Consumer<byte[]> consumer = client.newConsumer()
     .topic(topic)
     .subscriptionName("test")
     .autoAckOldestChunkedMessageOnQueueFull(true)
     .maxPendingChunkedMessage(100)
     .expireTimeOfIncompleteChunkedMessage(10, TimeUnit.MINUTES)
     .subscribe();
ConsumerConfiguration conf;
conf.setAutoAckOldestChunkedMessageOnQueueFull(true);
conf.setMaxPendingChunkedMessage(100);
Consumer consumer;
client.subscribe("my-topic", "my-sub", conf, consumer);

(Go: 준비 중이에요...)

consumer = client.subscribe(topic, "my-subscription",
                 max_pending_chunked_message=10,
                 auto_ack_oldest_chunked_message_on_queue_full=False
                 )

메시지 리스너를 가진 컨슈머 만들기 (Create a consumer with a message listener)

블로킹 호출로 루프를 실행하는 대신, 받은 각 메시지에 대해 호출되는 메시지 리스너를 사용하는 이벤트 기반 스타일을 쓸 수 있어요.

Consumer<String> consumer = pulsarClient.newConsumer(Schema.STRING)
                      .topic("persistent://my-property/my-ns/my-topic")
                      .subscriptionName("my-subscription")
                      .messageListener((c, m) -> {
                          try {
                              c.acknowledge(m);
                          } catch (Exception e) {
                              Assert.fail("Failed to acknowledge", e);
                          }
                      })
                      .subscribe();

이 예시는 가장 이른(earliest) 오프셋에서 구독을 시작하고 100개의 메시지를 소비해요.

#include <pulsar/Client.h>
#include <atomic>
#include <thread>
using namespace pulsar;

std::atomic<uint32_t> messagesReceived;

void handleAckComplete(Result res) {
    std::cout << "Ack res: " << res << std::endl;
}

void listener(Consumer consumer, const Message & msg) {
    std::cout << "Got message " << msg << " with content '" << msg.getDataAsString() << "'" << std::endl;
    messagesReceived++;
    consumer.acknowledgeAsync(msg.getMessageId(), handleAckComplete);
}

int main() {
    Client client("pulsar://localhost:6650");
    Consumer consumer;
    ConsumerConfiguration config;
    config.setMessageListener(listener);
    config.setSubscriptionInitialPosition(InitialPositionEarliest);
    Result result = client.subscribe("persistent://public/default/my-topic", "consumer-1", config, consumer);
    if (result != ResultOk) {
        std::cout << "Failed to subscribe: " << result << std::endl;
        return -1;
    }
    // wait for 100 messages to be consumed
    while (messagesReceived < 100) {
        std::this_thread::sleep_for(std::chrono::milliseconds(100));
    }
    std::cout << "Finished consuming asynchronously!" << std::endl;
    client.close();
    return 0;
}
import (
    "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()

    // we can listen this channel
    channel := make(chan pulsar.ConsumerMessage, 100)

    options := pulsar.ConsumerOptions{
        Topic:            "topic-1",
        SubscriptionName: "my-subscription",
        Type:             pulsar.Shared,
        // fill `MessageChannel` field will create a listener
        MessageChannel: channel,
    }

    consumer, err := client.Subscribe(options)
    if err != nil {
        log.Fatal(err)
    }

    defer consumer.Close()

    // Receive messages from channel. The channel returns a struct `ConsumerMessage` which contains message and the consumer from where
    // the message was received. It's not necessary here since we have 1 single consumer, but the channel could be
    // shared across multiple consumers as well
    for cm := range channel {
        consumer := cm.Consumer
        msg := cm.Message
        fmt.Printf("Consumer %s received a message, msgId: %v, content: '%s'\n",
            consumer.Name(), msg.ID(), string(msg.Payload()))

        consumer.Ack(msg)
    }
}
def my_listener(consumer, message):
    try:
        print("Received message: '{}' id='{}'".format(
            message.data(),
            message.message_id()
        ))
        consumer.acknowledge(message)
    except Exception as e:
        consumer.negative_acknowledge(message)
        print("Error processing message:", e)

consumer = client.subscribe(
    'persistent://my-property/my-ns/my-topic',
    'my-subscription',
    message_listener=my_listener
)

메시지 인터셉트 (Intercept messages)

ConsumerInterceptor는 컨슈머가 받은 메시지를 가로채고, 필요하면 변경(mutate)할 수 있어요. 인터페이스에는 6개의 주요 이벤트가 있어요.

  • beforeConsumereceive() 또는 receiveAsync()가 메시지를 반환하기 전에 트리거돼요. 이 이벤트 안에서 메시지를 수정할 수 있어요.
  • onAcknowledge는 컨슈머가 ack를 브로커로 보내기 전에 트리거돼요.
  • onAcknowledgeCumulative는 컨슈머가 누적 ack를 브로커로 보내기 전에 트리거돼요.
  • onNegativeAcksSend는 부정 ack로 인한 재전송이 발생할 때 트리거돼요.
  • onAckTimeoutSend는 ack 타임아웃으로 인한 재전송이 발생할 때 트리거돼요.
  • onPartitionsChange는 (파티셔닝된) 토픽의 파티션이 변경될 때 트리거돼요.

메시지를 인터셉트하려면 아래처럼 Consumer를 만들 때 하나 또는 여러 ConsumerInterceptor를 추가하면 돼요.

Consumer<String> consumer = client.newConsumer()
     .topic("my-topic")
     .subscriptionName("my-subscription")
     .intercept(new ConsumerInterceptor<String> {
           @Override
           public Message<String> beforeConsume(Consumer<String> consumer, Message<String> message) {
               // user-defined processing logic
           }

           @Override
           public void onAcknowledge(Consumer<String> consumer, MessageId messageId, Throwable cause) {
               // user-defined processing logic
           }

           @Override
           public void onAcknowledgeCumulative(Consumer<String> consumer, MessageId messageId, Throwable cause) {
               // user-defined processing logic
           }

           @Override
           public void onNegativeAcksSend(Consumer<String> consumer, Set<MessageId> messageIds) {
               // user-defined processing logic
           }

           @Override
           public void onAckTimeoutSend(Consumer<String> consumer, Set<MessageId> messageIds) {
               // user-defined processing logic
           }

           @Override
           public void onPartitionsChange(String topicName, int partitions) {
               // user-defined processing logic
           }
     })
     .subscribe();

여러 인터셉터를 사용한다면, intercept 메서드에 전달된 순서대로 적용돼요.

더 알아보기 (Learn more)