All Products
Search
Document Center

ApsaraMQ for Kafka:Kirim dan berlangganan pesan dengan Go SDK

Last Updated:Jul 11, 2026

Gunakan Go SDK untuk terhubung ke instans Message Queue for Apache Kafka melalui titik akhir (endpoint), lalu kirim dan berlangganan pesan.

Sebelum memulai

Go telah diinstal. Untuk informasi lebih lanjut, lihat Install Go.

Catatan

Demo kafka-confluent-go-demo tidak mendukung Windows.

Persiapkan konfigurasi

  1. Opsi: Unduh sertifikat root SSL. Sertifikat ini diperlukan saat menghubungkan ke instans menggunakan endpoint SSL.

  2. Buka repositori aliware-kafka-demos, klik ikon code, lalu pilih Download ZIP dari daftar drop-down untuk mengunduh proyek demo. Setelah pengunduhan selesai, ekstrak file tersebut.

  3. Di dalam proyek demo yang telah diekstrak, temukan folder kafka-confluent-go-demo dan unggah ke direktori /home sistem Linux Anda.

  4. Masuk ke sistem Linux Anda, buka direktori /home/kafka-confluent-go-demo, lalu ubah file konfigurasi conf/kafka.json.

    {
      "topic": "XXX",
      "group.id": "XXX",
      "bootstrap.servers" : "XXX:XX,XXX:XX,XXX:XX",
      "security.protocol" : "plaintext",
      "sasl.mechanism" : "XXX",
      "sasl.username" : "XXX",
      "sasl.password" : "XXX"
    }

    Parameter

    Deskripsi

    topic

    Nama topik di instans Anda. Anda dapat menemukan nama topik tersebut di halaman Topics pada Message Queue for Apache Kafka console.

    group.id

    ID group konsumen di instans Anda. Anda dapat menemukan ID group tersebut di halaman Groups pada Message Queue for Apache Kafka console.

    Catatan

    Parameter ini opsional saat mengirim pesan menggunakan producer.go, tetapi wajib diisi saat berlangganan pesan menggunakan consumer.go.

    bootstrap.servers

    Alamat IP dan port titik akhir (endpoint). Informasi ini dapat ditemukan di bagian Endpoint Information pada halaman Instance Details di Message Queue for Apache Kafka console.

    security.protocol

    Protokol untuk otentikasi pengguna Simple Authentication and Security Layer (SASL). Nilai default-nya adalah plaintext. Nilainya bergantung pada jenis endpoint:

    • Endpoint default: plaintext.

    • Endpoint SSL: sasl_ssl.

    • Endpoint SASL: sasl_plaintext.

    sasl.mechanism

    Mekanisme keamanan untuk mengirim dan berlangganan pesan. Nilainya bergantung pada jenis endpoint:

    • Endpoint default: Parameter ini tidak diperlukan.

    • Endpoint SSL: PLAIN.

    • Endpoint SASL: Untuk mekanisme PLAIN, atur parameter ini ke PLAIN. Untuk mekanisme SCRAM, atur parameter ini ke SCRAM-SHA-256.

    sasl.username

    Username SASL. Parameter ini wajib diisi jika Anda menggunakan endpoint SSL atau endpoint SASL.

    Catatan
    • Jika fitur ACL tidak diaktifkan untuk instans ApsaraMQ for Kafka, Anda dapat memperoleh username dan password pengguna SASL dari parameter Username dan Password di bagian Configuration Information pada halaman Instance Details di ApsaraMQ for Kafka console.

    • Jika fitur ACL diaktifkan untuk instans ApsaraMQ for Kafka, pastikan pengguna SASL telah diberi izin untuk mengirim dan menerima pesan menggunakan instans tersebut. Untuk informasi selengkapnya, lihat Grant permissions to SASL users.

    sasl.password

    Password SASL. Parameter ini wajib diisi jika Anda menggunakan endpoint SSL atau endpoint SASL.

Kirim pesan

Jalankan perintah berikut untuk mengirim pesan menggunakan file producer.go.

go run -mod=vendor producer/producer.go

Contoh file producer.go:

package main

import (
    "encoding/json"
    "fmt"
    "github.com/confluentinc/confluent-kafka-go/kafka"
    "os"
    "path/filepath"
)

type KafkaConfig struct {
    Topic      string `json:"topic"`
    GroupId    string `json:"group.id"`
    BootstrapServers    string `json:"bootstrap.servers"`
    SecurityProtocol string `json:"security.protocol"`
    SslCaLocation string `json:"ssl.ca.location"`
    SaslMechanism string `json:"sasl.mechanism"`
    SaslUsername string `json:"sasl.username"`
    SaslPassword string `json:"sasl.password"`
}

// Parameter config harus berupa pointer ke struct; jika tidak, fungsi akan panic.
func loadJsonConfig() *KafkaConfig {
    workPath, err := os.Getwd()
    if err != nil {
        panic(err)
    }
    configPath := filepath.Join(workPath, "conf")
    fullPath := filepath.Join(configPath, "kafka.json")
    file, err := os.Open(fullPath);
    if (err != nil) {
        msg := fmt.Sprintf("Tidak dapat memuat konfigurasi di %s. Error: %v", fullPath, err)
        panic(msg)
    }

    defer file.Close()

    decoder := json.NewDecoder(file)
    var config = &KafkaConfig{}
    err = decoder.Decode(config);
    if (err != nil) {
        msg := fmt.Sprintf("Gagal decode JSON untuk file konfigurasi di %s. Error: %v", fullPath, err)
        panic(msg)
    }
    json.Marshal(config)
    return  config
}

func doInitProducer(cfg *KafkaConfig) *kafka.Producer {
    fmt.Print("inisialisasi produsen kafka, proses ini mungkin memerlukan beberapa detik untuk membuat koneksi\n")
    // argumen umum
    var kafkaconf = &kafka.ConfigMap{
        "api.version.request": "true",
        "message.max.bytes": 1000000,
        "linger.ms": 10,
        "retries": 30,
        "retry.backoff.ms": 1000,
        "acks": "1"}
    kafkaconf.SetKey("bootstrap.servers", cfg.BootstrapServers)

    switch cfg.SecurityProtocol {
        case "plaintext" :
            kafkaconf.SetKey("security.protocol", "plaintext");
        case "sasl_ssl":
            kafkaconf.SetKey("security.protocol", "sasl_ssl");
            kafkaconf.SetKey("ssl.ca.location", "conf/ca-cert.pem");
            kafkaconf.SetKey("sasl.username", cfg.SaslUsername);
            kafkaconf.SetKey("sasl.password", cfg.SaslPassword);
            kafkaconf.SetKey("sasl.mechanism", cfg.SaslMechanism);
            kafkaconf.SetKey("enable.ssl.certificate.verification", "false");
            kafkaconf.SetKey("ssl.endpoint.identification.algorithm", "None")
    case "sasl_plaintext":
            kafkaconf.SetKey("sasl.mechanism", "PLAIN")
            kafkaconf.SetKey("security.protocol", "sasl_plaintext");
            kafkaconf.SetKey("sasl.username", cfg.SaslUsername);
            kafkaconf.SetKey("sasl.password", cfg.SaslPassword);
            kafkaconf.SetKey("sasl.mechanism", cfg.SaslMechanism)
    default:
            panic(kafka.NewError(kafka.ErrUnknownProtocol, "protokol tidak dikenal", true))
    }

    producer, err := kafka.NewProducer(kafkaconf)
    if err != nil {
        panic(err)
    }
    fmt.Print("inisialisasi produsen kafka berhasil\n")
    return producer;
}

func main() {
    // Pilih protokol yang sesuai
    // 9092 untuk PLAINTEXT
    // 9093 untuk SASL_SSL, perlu menyediakan sasl.username dan sasl.password
    // 9094 untuk SASL_PLAINTEXT, perlu menyediakan sasl.username dan sasl.password
    cfg := loadJsonConfig();
    producer := doInitProducer(cfg)

    defer producer.Close()

    // Penanganan laporan pengiriman untuk pesan yang diproduksi
    go func() {
        for e := range producer.Events() {
            switch ev := e.(type) {
            case *kafka.Message:
                if ev.TopicPartition.Error != nil {
                    fmt.Printf("Pengiriman gagal: %v\n", ev.TopicPartition)
                } else {
                    fmt.Printf("Pesan berhasil dikirim ke %v\n", ev.TopicPartition)
                }
            }
        }
    }()

    // Produksi pesan ke topik (secara asinkron)
    topic := cfg.Topic
    for _, word := range []string{"Welcome", "to", "the", "Confluent", "Kafka", "Golang", "client"} {
        producer.Produce(&kafka.Message{
            TopicPartition: kafka.TopicPartition{Topic: &topic, Partition: kafka.PartitionAny},
            Value:          []byte(word),
        }, nil)
    }

    // Tunggu hingga semua pesan terkirim sebelum menutup
    producer.Flush(15 * 1000)
}
                        

Berlangganan pesan

Jalankan perintah berikut untuk berlangganan pesan menggunakan file consumer.go.

go run -mod=vendor consumer/consumer.go

Contoh file consumer.go:

package main

import (
    "encoding/json"
    "fmt"
    "github.com/confluentinc/confluent-kafka-go/kafka"
    "os"
    "path/filepath"
)
type KafkaConfig struct {
    Topic      string `json:"topic"`
    GroupId    string `json:"group.id"`
    BootstrapServers    string `json:"bootstrap.servers"`
    SecurityProtocol string `json:"security.protocol"`
    SaslMechanism string `json:"sasl.mechanism"`
    SaslUsername string `json:"sasl.username"`
    SaslPassword string `json:"sasl.password"`
}

// Parameter config harus berupa pointer ke struct; jika tidak, fungsi akan panic.
func loadJsonConfig() *KafkaConfig {
    workPath, err := os.Getwd()
    if err != nil {
        panic(err)
    }
    configPath := filepath.Join(workPath, "conf")
    fullPath := filepath.Join(configPath, "kafka.json")
    file, err := os.Open(fullPath);
    if (err != nil) {
        msg := fmt.Sprintf("Tidak dapat memuat konfigurasi di %s. Error: %v", fullPath, err)
        panic(msg)
    }

    defer file.Close()

    decoder := json.NewDecoder(file)
    var config = &KafkaConfig{}
    err = decoder.Decode(config);
    if (err != nil) {
        msg := fmt.Sprintf("Gagal decode JSON untuk file konfigurasi di %s. Error: %v", fullPath, err)
        panic(msg)
    }
    json.Marshal(config)
    return  config
}


func doInitConsumer(cfg *KafkaConfig) *kafka.Consumer {
    fmt.Print("inisialisasi konsumen kafka, proses ini mungkin memerlukan beberapa detik untuk membuat koneksi\n")
    // argumen umum
    var kafkaconf = &kafka.ConfigMap{
        "api.version.request": "true",
        "auto.offset.reset": "latest",
        "heartbeat.interval.ms": 3000,
        "session.timeout.ms": 30000,
        "max.poll.interval.ms": 120000,
        "fetch.max.bytes": 1024000,
        "max.partition.fetch.bytes": 256000}
    kafkaconf.SetKey("bootstrap.servers", cfg.BootstrapServers);
    kafkaconf.SetKey("group.id", cfg.GroupId)

    switch cfg.SecurityProtocol {
    case "plaintext" :
        kafkaconf.SetKey("security.protocol", "plaintext");
    case "sasl_ssl":
        kafkaconf.SetKey("security.protocol", "sasl_ssl");
        kafkaconf.SetKey("ssl.ca.location", "./conf/ca-cert.pem");
        kafkaconf.SetKey("sasl.username", cfg.SaslUsername);
        kafkaconf.SetKey("sasl.password", cfg.SaslPassword);
        kafkaconf.SetKey("sasl.mechanism", cfg.SaslMechanism);
        kafkaconf.SetKey("ssl.endpoint.identification.algorithm", "None");
        kafkaconf.SetKey("enable.ssl.certificate.verification", "false")
    case "sasl_plaintext":
        kafkaconf.SetKey("security.protocol", "sasl_plaintext");
        kafkaconf.SetKey("sasl.username", cfg.SaslUsername);
        kafkaconf.SetKey("sasl.password", cfg.SaslPassword);
        kafkaconf.SetKey("sasl.mechanism", cfg.SaslMechanism)

    default:
        panic(kafka.NewError(kafka.ErrUnknownProtocol, "protokol tidak dikenal", true))
    }

    consumer, err := kafka.NewConsumer(kafkaconf)
    if err != nil {
        panic(err)
    }
    fmt.Print("inisialisasi konsumen kafka berhasil\n")
    return consumer;
}

func main() {

    // Pilih protokol yang sesuai
    // 9092 untuk PLAINTEXT
    // 9093 untuk SASL_SSL, perlu menyediakan sasl.username dan sasl.password
    // 9094 untuk SASL_PLAINTEXT, perlu menyediakan sasl.username dan sasl.password
    cfg := loadJsonConfig();
    consumer := doInitConsumer(cfg)

    consumer.SubscribeTopics([]string{cfg.Topic}, nil)

    for {
        msg, err := consumer.ReadMessage(-1)
        if err == nil {
            fmt.Printf("Pesan pada %s: %s\n", msg.TopicPartition, string(msg.Value))
        } else {
            // Klien secara otomatis akan mencoba memulihkan diri dari semua error.
            fmt.Printf("Error konsumen: %v (%v)\n", err, msg)
        }
    }

    consumer.Close()
}