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.
Demo kafka-confluent-go-demo tidak mendukung Windows.
Persiapkan konfigurasi
-
Opsi: Unduh sertifikat root SSL. Sertifikat ini diperlukan saat menghubungkan ke instans menggunakan endpoint SSL.
-
Buka repositori aliware-kafka-demos, klik ikon
, lalu pilih Download ZIP dari daftar drop-down untuk mengunduh proyek demo. Setelah pengunduhan selesai, ekstrak file tersebut. -
Di dalam proyek demo yang telah diekstrak, temukan folder kafka-confluent-go-demo dan unggah ke direktori /home sistem Linux Anda.
-
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.
CatatanParameter 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.
CatatanJika 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()
}