컨슈머 사용하기
컨슈머 사용하기 (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();
위 예시에서 컨슈머는 토픽 이름 패턴에 일치하는 영속 토픽을 구독해요. 컨슈머가 토픽 이름 패턴에 일치하는 모든 영속 및 비영속 토픽을 구독하려면 subscriptionTopicsMode를 RegexSubscriptionMode.AllTopics로 설정하세요.
Pattern pattern = Pattern.compile("public/default/.*");
pulsarClient.newConsumer()
.subscriptionName("my-sub")
.topicsPattern(pattern)
.subscriptionTopicsMode(RegexSubscriptionMode.AllTopics)
.subscribe();
참고: 기본적으로 컨슈머의
subscriptionTopicsMode는PersistentOnly예요.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 |
negativeAckRedeliveryBackoff는consumer.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개의 주요 이벤트가 있어요.
beforeConsume은receive()또는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)
- 프로듀서 사용하기 — 프로듀서로 메시지를 보내는 방법을 살펴봐요.
- 리더 사용하기 — 리더로 메시지를 읽는 방법을 알아봐요.
- 서브스크립션 개념 — 서브스크립션 타입의 개념을 살펴봐요.
- 참조: 아키텍처 개념 — ack 등 메시징 보장을 확인해요.
- Java 클라이언트 — Java 클라이언트 사용법을 확인해요.