Вопрос проверяет понимание особенностей работы с Apache Kafka на языке Go, включая использование библиотек, обработку партиций и управление потребителями.
Apache Kafka — это распределенная платформа для потоковой передачи данных. При разработке на Go необходимо учитывать особенности языка и доступные библиотеки для взаимодействия с Kafka. Основные библиотеки: Sarama (чистый Go) и Confluent Kafka Go (обертка над librdkafka на C).
package main
import (
"fmt"
"github.com/Shopify/sarama"
)
func main() {
config := sarama.NewConfig()
config.Producer.RequiredAcks = sarama.WaitForAll
config.Producer.Retry.Max = 5
producer, err := sarama.NewSyncProducer([]string{"localhost:9092"}, config)
if err != nil {
panic(err)
}
defer producer.Close()
msg := &sarama.ProducerMessage{
Topic: "test",
Key: sarama.StringEncoder("key"),
Value: sarama.StringEncoder("Hello Kafka from Go"),
}
partition, offset, err := producer.SendMessage(msg)
if err != nil {
panic(err)
}
fmt.Printf("Message sent to partition %d at offset %d\n", partition, offset)
}package main
import (
"fmt"
"github.com/Shopify/sarama"
)
func main() {
config := sarama.NewConfig()
config.Consumer.Return.Errors = true
consumer, err := sarama.NewConsumer([]string{"localhost:9092"}, config)
if err != nil {
panic(err)
}
defer consumer.Close()
partitionConsumer, err := consumer.ConsumePartition("test", 0, sarama.OffsetNewest)
if err != nil {
panic(err)
}
defer partitionConsumer.Close()
for msg := range partitionConsumer.Messages() {
fmt.Printf("Received message: %s\n", string(msg.Value))
}
}Работа с Kafka на Go требует понимания асинхронной обработки и управления состоянием потребителей. Используйте Sarama для простых проектов или Confluent Kafka Go для высокой производительности. Правильная обработка ошибок и настройка параметров гарантируют надежную передачу данных.