Go 클라이언트 사용

Go 클라이언트 사용 (Go client use)

Go 클라이언트로 프로듀서, 컨슈머, 리더를 만들어 메시지를 보내고 받는 방법을 살펴볼게요. 특히 Prometheus 메트릭을 노출하는 간단한 HTTP 서버 예시로, 운영 환경에서 클라이언트 상태를 어떻게 관찰하는지도 함께 확인할 수 있어요.

출처: 문서

본문

프로듀서 만들기 (Create a producer)

ProducerOptions 객체를 사용해 Go 프로듀서를 구성할 수 있어요. 다음은 예시예요.

producer, err := client.CreateProducer(pulsar.ProducerOptions{
    Topic: "my-topic",
})

if err != nil {
    log.Fatal(err)
}

_, err = producer.Send(context.Background(), &pulsar.ProducerMessage{
    Payload: []byte("hello"),
})

defer producer.Close()

if err != nil {
    fmt.Println("Failed to publish message", err)
}
fmt.Println("Published message")

Producer 인터페이스의 모든 사용 가능한 메서드는 여기를 참고하세요.

모니터링 (Monitor)

Pulsar Go 클라이언트는 Prometheus를 사용해 클라이언트 메트릭을 등록해요. 이 섹션에서는 Prometheus 메트릭을 HTTP로 노출하는 간단한 Pulsar 프로듀서 애플리케이션을 만드는 방법을 보여드릴게요.

  • 간단한 프로듀서 애플리케이션 작성:
// Create a Pulsar client
client, err := pulsar.NewClient(pulsar.ClientOptions{
    URL: "pulsar://localhost:6650",
})
if err != nil {
    log.Fatal(err)
}

defer client.Close()

// Start a separate goroutine for Prometheus metrics
// In this case, Prometheus metrics can be accessed via http://localhost:2112/metrics
go func() {
    prometheusPort := 2112
    log.Printf("Starting Prometheus metrics at http://localhost:%v/metrics\n", prometheusPort)
    http.Handle("/metrics", promhttp.Handler())
    err = http.ListenAndServe(":"+strconv.Itoa(prometheusPort), nil)
    if err != nil {
        log.Fatal(err)
    }
}()

// Create a producer
producer, err := client.CreateProducer(pulsar.ProducerOptions{
    Topic: "topic-1",
})
if err != nil {
    log.Fatal(err)
}

defer producer.Close()

ctx := context.Background()

// Write your business logic here
// In this case, you build a simple Web server. You can produce messages by requesting http://localhost:8082/produce
webPort := 8082
http.HandleFunc("/produce", func(w http.ResponseWriter, r *http.Request) {
    msgId, err := producer.Send(ctx, &pulsar.ProducerMessage{
        Payload: []byte(fmt.Sprintf("hello world")),
    })
    if err != nil {
        log.Fatal(err)
    } else {
        log.Printf("Published message: %v", msgId)
        fmt.Fprintf(w, "Published message: %v", msgId)
    }
})

err = http.ListenAndServe(":"+strconv.Itoa(webPort), nil)
if err != nil {
    log.Fatal(err)
}
  • 애플리케이션에서 메트릭을 스크랩하려면 구성 파일(prometheus.yml)을 사용해 로컬에서 실행 중인 Prometheus 인스턴스를 구성해요.
scrape_configs:
- job_name: pulsar-client-go-metrics
  scrape_interval: 10s
  static_configs:
  - targets:
  - localhost:2112

컨슈머 만들기 (Create a consumer)

Pulsar 컨슈머는 하나 이상의 Pulsar 토픽을 구독하고, 그 토픽에서 생산된 들어오는 메시지를 기다려요. ConsumerOptions 객체를 사용해 Go 컨슈머를 구성할 수 있어요. 채널을 사용하는 기본 예시는 다음과 같아요.

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

for i := 0; i < 10; i++ {
    // may block here
    msg, err := consumer.Receive(context.Background())
    if err != nil {
        log.Fatal(err)
    }

    fmt.Printf("Received message msgId: %#v -- content: '%s'\n",
        msg.ID(), string(msg.Payload()))

    consumer.Ack(msg)
}

if err := consumer.Unsubscribe(); err != nil {
    log.Fatal(err)
}

Consumer 인터페이스의 모든 사용 가능한 메서드는 여기를 참고하세요.

단일 토픽 컨슈머 만들기 (Create a single-topic consumer)

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

defer client.Close()

consumer, err := client.Subscribe(pulsar.ConsumerOptions{
    // fill `Topic` field will create a single-topic consumer
    Topic:            "topic-1",
    SubscriptionName: "my-sub",
    Type:             pulsar.Shared,
})
if err != nil {
    log.Fatal(err)
}

defer consumer.Close()

정규식 토픽 컨슈머 만들기 (Create a regex-topic consumer)

client, err := pulsar.NewClient(pulsar.ClientOptions{
    URL: "pulsar://localhost:6650",
})
defer client.Close()

topicsPattern := "persistent://public/default/topic.*"
opts := pulsar.ConsumerOptions{
    // fill `TopicsPattern` field will create a regex consumer
    TopicsPattern:    topicsPattern,
    SubscriptionName: "regex-sub",
}

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

모니터링 (Monitor)

이 섹션에서는 Prometheus 메트릭을 HTTP로 노출하는 간단한 Pulsar 컨슈머 애플리케이션을 만드는 방법을 보여드릴게요.

  • 간단한 컨슈머 애플리케이션 작성:
// Create a Pulsar client
client, err := pulsar.NewClient(pulsar.ClientOptions{
    URL: "pulsar://localhost:6650",
})
if err != nil {
    log.Fatal(err)
}

defer client.Close()

// Start a separate goroutine for Prometheus metrics
// In this case, Prometheus metrics can be accessed via http://localhost:2112/metrics
go func() {
    prometheusPort := 2112
    log.Printf("Starting Prometheus metrics at http://localhost:%v/metrics\n", prometheusPort)
    http.Handle("/metrics", promhttp.Handler())
    err = http.ListenAndServe(":"+strconv.Itoa(prometheusPort), nil)
    if err != nil {
        log.Fatal(err)
    }
}()

// Create a consumer
consumer, err := client.Subscribe(pulsar.ConsumerOptions{
    Topic:            "topic-1",
    SubscriptionName: "sub-1",
    Type:             pulsar.Shared,
})
if err != nil {
    log.Fatal(err)
}

defer consumer.Close()

ctx := context.Background()

// Write your business logic here
// In this case, you build a simple Web server. You can consume messages by requesting http://localhost:8083/consume
webPort := 8083
http.HandleFunc("/consume", func(w http.ResponseWriter, r *http.Request) {
    msg, err := consumer.Receive(ctx)
    if err != nil {
        log.Fatal(err)
    } else {
        log.Printf("Received message msgId: %v -- content: '%s'\n", msg.ID(), string(msg.Payload()))
        fmt.Fprintf(w, "Received message msgId: %v -- content: '%s'\n", msg.ID(), string(msg.Payload()))
        consumer.Ack(msg)
    }
})

err = http.ListenAndServe(":"+strconv.Itoa(webPort), nil)
if err != nil {
    log.Fatal(err)
}
  • 애플리케이션에서 메트릭을 스크랩하려면 구성 파일(prometheus.yml)을 사용해 로컬에서 실행 중인 Prometheus 인스턴스를 구성해요.
scrape_configs:
- job_name: pulsar-client-go-metrics
  scrape_interval: 10s
  static_configs:
  - targets:
  - localhost: 2112

리더 만들기 (Create a reader)

Pulsar 리더는 Pulsar 토픽의 메시지를 처리해요. 리더는 컨슈머와 다른데, 리더는 스트림에서 어떤 메시지부터 시작할지 명시적으로 지정해야 하기 때문이에요(반면 컨슈머는 자동으로 가장 최근의 ack되지 않은 메시지부터 시작해요). ReaderOptions 객체를 사용해 Go 리더를 구성할 수 있어요. 다음은 예시예요.

reader, err := client.CreateReader(pulsar.ReaderOptions{
    Topic:          "topic-1",
    StartMessageID: pulsar.EarliestMessageID(),
})
if err != nil {
    log.Fatal(err)
}
defer reader.Close()

Reader 인터페이스의 모든 사용 가능한 메서드는 여기를 참고하세요.

더 알아보기 (Learn more)