Pulsar를 메시지 큐로 사용하기

Pulsar를 메시지 큐로 사용하기

시스템의 어떤 컴포넌트가 느려지거나 실패하더라도 작업 객체를 반드시 처리해야 하는 경우가 있어요. 그럴 때 메시지 큐가 개입해서 올바른 순서로 처리되지 않은 데이터를 보존해 주면 안심할 수 있죠. Pulsar는 영구 메시지 저장을 염두에 두고 만들어진 데다 소비자 간 자동 부하 분산을 지원해서 메시지 큐로 아주 좋은 선택이에요. 이 글에서는 Pulsar 토픽을 메시지 큐처럼 사용하는 방법을 정리해 드릴게요.

출처: 문서

본문

메시지 큐는 많은 대규모 데이터 아키텍처의 필수 구성 요소예요. 시스템을 통과하는 모든 작업 객체가 시스템 구성 요소의 느림이나 완전한 실패에도 불구하고 반드시 처리되어야 한다면, 메시지 큐가 개입해서 필요한 조치가 취해질 때까지 처리되지 않은 데이터를 올바른 순서로 보존해 줄 가능성이 높아요.

Pulsar는 메시지 큐로 아주 좋은 선택이에요. 왜냐하면:

  • 영구 메시지 저장을 염두에 두고 만들어졌어요.
  • 토픽의 메시지에 대해 소비자 간 자동 부하 분산(원한다면 커스텀 부하 분산도)을 제공해요.

팁 같은 Pulsar 설치를 실시간 메시지 버스로도, 메시지 큐로도 사용할 수 있어요 (또는 둘 중 하나만). 일부 토픽은 실시간 용도로, 다른 토픽은 메시지 큐 용도로 따로 떼어 놓을 수 있어요 (또는 원하면 특정 네임스페이스를 각 용도에 사용해도 돼요).

클라이언트 구성 변경

Pulsar 토픽을 메시지 큐로 사용하려면 그 토픽의 수신기 부하를 여러 소비자에 분산해야 해요 (최적의 소비자 수는 부하에 따라 달라져요).

각 소비자는 shared 구독을 설정하고 다른 소비자와 같은 구독 이름을 사용해야 해요 (그렇지 않으면 구독이 공유되지 않고 소비자들이 처리 앙상블로 동작할 수 없어요).

메시지 전달을 소비자 간에 엄격하게 제어하고 싶다면 소비자의 수신기 큐(receiver queue) 크기를 아주 낮게 설정해요 (필요하면 0까지도 가능해요). 각 소비자는 한 번에 가져오려는 메시지 수를 결정하는 수신기 큐를 가져요. 예를 들어 수신기 큐가 1000(기본값)이면, 연결 시 소비자가 토픽 백로그에서 1000개 메시지를 처리하려 시도한다는 뜻이에요. 수신기 큐를 0으로 설정하면 본질적으로 각 소비자가 한 번에 한 가지 작업만 하도록 보장하는 거예요.

팁 파티션 토픽 소비자의 수신기 큐 크기는 다음 두 값 중 최솟값을 채택해요:

  • receiver_queue_size
  • max_total_receiver_queue_size_across_partitions/NumPartitions

예시

여기에 shared 구독을 사용하는 예시가 있어요.

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

String SERVICE_URL = "pulsar://localhost:6650";
String TOPIC = "persistent://public/default/mq-topic-1";
String subscription = "sub-1";

PulsarClient client = PulsarClient.builder()
        .serviceUrl(SERVICE_URL)
        .build();

Consumer consumer = client.newConsumer()
        .topic(TOPIC)
        .subscriptionName(subscription)
        .subscriptionType(SubscriptionType.Shared)
        // If you'd like to restrict the receiver queue size
        .receiverQueueSize(10)
        .subscribe();
from pulsar import Client, ConsumerType

SERVICE_URL = "pulsar://localhost:6650"
TOPIC = "persistent://public/default/mq-topic-1"
SUBSCRIPTION = "sub-1"

client = Client(SERVICE_URL)
consumer = client.subscribe(
    TOPIC,
    SUBSCRIPTION,
    # If you'd like to restrict the receiver queue size
    receiver_queue_size=10,
    consumer_type=ConsumerType.Shared)
#include <pulsar/Client.h>

std::string serviceUrl = "pulsar://localhost:6650";
std::string topic = "persistent://public/default/mq-topic-1";
std::string subscription = "sub-1";

Client client(serviceUrl);
ConsumerConfiguration consumerConfig;
consumerConfig.setConsumerType(ConsumerType.ConsumerShared);
// If you'd like to restrict the receiver queue size
consumerConfig.setReceiverQueueSize(10);

Consumer consumer;
Result result = client.subscribe(topic, subscription, consumerConfig, consumer);
import "github.com/apache/pulsar-client-go/pulsar"

client, err := pulsar.NewClient(pulsar.ClientOptions{
    URL: "pulsar://localhost:6650",
})
if err != nil {
    log.Fatal(err)
}

consumer, err := client.Subscribe(pulsar.ConsumerOptions{
    Topic:             "persistent://public/default/mq-topic-1",
    SubscriptionName:  "sub-1",
    Type:              pulsar.Shared,
    ReceiverQueueSize: 10, // If you'd like to restrict the receiver queue size
})
if err != nil {
    log.Fatal(err)
}

더 알아보기 (Learn more)