Producer на Go: публикуем события в Kafka (franz-go)
Producer на Go: публикуем события в Kafka (franz-go)
franz-go - современная Go-библиотека для Kafka. Написана с нуля в 2021, поддерживает все Kafka API начиная с версии 0.8, активно развивается. В отличие от sarama, у неё понятный API и нет накопленного технического долга.
go get github.com/twmb/franz-go/pkg/kgo
Создание клиента
package main
import (
"context"
"fmt"
"log"
"github.com/twmb/franz-go/pkg/kgo"
)
func main() {
client, err := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
kgo.DefaultProduceTopic("lesson-events"),
)
if err != nil {
log.Fatal(err)
}
defer client.Close()
}
SeedBrokers - bootstrap-серверы. Клиент подключится к ним, получит метаданные кластера (список всех брокеров, топиков, партиций) и дальше работает самостоятельно.
Синхронная отправка
record := &kgo.Record{
Topic: "lesson-events",
Key: []byte("user-42"), // partition key
Value: []byte(`{"event":"lesson.completed","lessonId":"go-01"}`),
Headers: []kgo.RecordHeader{
{Key: "content-type", Value: []byte("application/json")},
{Key: "version", Value: []byte("1")},
},
}
results := client.ProduceSync(ctx, record)
if err := results.FirstErr(); err != nil {
return fmt.Errorf("produce: %w", err)
}
ProduceSync блокирует до получения подтверждения от брокера. Удобно для простых случаев, но медленно при высоком throughput - каждая отправка ждёт ack.
Асинхронная отправка
var wg sync.WaitGroup
for _, event := range events {
wg.Add(1)
record := &kgo.Record{
Topic: "lesson-events",
Key: []byte(event.UserID),
Value: marshal(event),
}
client.Produce(ctx, record, func(r *kgo.Record, err error) {
defer wg.Done()
if err != nil {
log.Printf("produce error: %v", err)
}
})
}
wg.Wait()
Produce (без Sync) принимает callback и возвращает сразу. Клиент буферизует записи и отправляет батчами - это даёт на порядок больший throughput.
Partition Key и порядок
// Все события одного пользователя идут в одну партицию
record := &kgo.Record{
Key: []byte(userID), // partition key
Value: payload,
}
franz-go хеширует Key и выбирает партицию. Одинаковый ключ = одна и та же партиция = порядок гарантирован. Это критично, если нужна последовательность событий: lesson.started → lesson.completed.
Без ключа (nil) - round-robin по партициям, порядок не гарантирован.
Производительная конфигурация
client, err := kgo.NewClient(
kgo.SeedBrokers("localhost:9092"),
// Батчинг: накапливать до 5ms перед отправкой
kgo.ProducerLinger(5 * time.Millisecond),
// Максимальный размер батча
kgo.ProducerBatchMaxBytes(1_000_000), // 1MB
// Сжатие (snappy - хороший баланс CPU vs compression ratio)
kgo.ProducerBatchCompression(kgo.SnappyCompression()),
// Требуем подтверждения от всех in-sync реплик
kgo.RequiredAcks(kgo.AllISRAcks()),
// Таймаут на доставку записи
kgo.RecordDeliveryTimeout(10 * time.Second),
)
ProducerLinger - ключевой параметр для throughput. Без него каждый Produce улетает немедленно. С ним клиент ждёт 5ms, чтобы накопить батч. 5-10ms - стандартное значение.
Graceful Shutdown
func runProducer(ctx context.Context) error {
client, err := kgo.NewClient(kgo.SeedBrokers("localhost:9092"))
if err != nil {
return err
}
defer func() {
// Flush ждёт пока все буферизованные записи будут отправлены
if err := client.Flush(ctx); err != nil {
log.Printf("flush error: %v", err)
}
client.Close()
}()
// ... produce events ...
return nil
}
client.Close() без Flush может потерять записи из буфера. Всегда вызывай Flush перед Close. В HTTP-сервере вызывай Flush в shutdown hook.
Producer как зависимость (Clean Architecture)
// domain/events.go
type LessonCompletedEvent struct {
UserID string
LessonID string
At time.Time
}
// adapters/kafka_publisher.go
type KafkaPublisher struct {
client *kgo.Client
topic string
}
func NewKafkaPublisher(brokers []string, topic string) (*KafkaPublisher, error) {
client, err := kgo.NewClient(
kgo.SeedBrokers(brokers...),
kgo.ProducerLinger(5*time.Millisecond),
kgo.RequiredAcks(kgo.AllISRAcks()),
)
if err != nil {
return nil, err
}
return &KafkaPublisher{client: client, topic: topic}, nil
}
func (p *KafkaPublisher) PublishLessonCompleted(ctx context.Context, e LessonCompletedEvent) error {
payload, _ := json.Marshal(e)
r := &kgo.Record{
Topic: p.topic,
Key: []byte(e.UserID),
Value: payload,
}
return p.client.ProduceSync(ctx, r).FirstErr()
}
func (p *KafkaPublisher) Close(ctx context.Context) error {
if err := p.client.Flush(ctx); err != nil {
return err
}
p.client.Close()
return nil
}
Use case работает с интерфейсом EventPublisher, не зная о Kafka - принцип инверсии зависимостей.
Обработка ошибок
results := client.ProduceSync(ctx, record)
for _, result := range results {
if result.Err != nil {
// Retryable: network error, leader election
// Non-retryable: message too large, topic not found
if errors.Is(result.Err, context.DeadlineExceeded) {
// таймаут - попробовать позже
} else {
log.Printf("permanent error for record: %v", result.Err)
}
}
}
franz-go автоматически ретраит транзиентные ошибки (leader election, временная недоступность брокера). RecordDeliveryTimeout ограничивает суммарное время всех попыток.
PHP: публикация через rdkafka
<?php
$conf = new RdKafka\Conf();
$conf->set('metadata.broker.list', 'localhost:9092');
$conf->set('socket.timeout.ms', '3000');
$producer = new RdKafka\Producer($conf);
$topic = $producer->newTopic('lesson-events');
$payload = json_encode(['event' => 'lesson.completed', 'lessonId' => 'go-01']);
$topic->produce(RD_KAFKA_PARTITION_UA, 0, $payload, $userID); // ключ = $userID
$producer->flush(3000); // ждём подтверждения
RD_KAFKA_PARTITION_UA - unassigned: брокер выберет партицию по ключу.
Типичные ошибки
Не вызывать Flush при завершении - последние N записей из буфера теряются при остановке сервиса. Особенно опасно с ProducerLinger: буфер может накопить события за 5ms, а потом сервис упадёт.
RequiredAcks(kgo.NoAck()) - отправка без подтверждения. Максимальный throughput, нулевая надёжность. Для продакшна минимум kgo.LeaderAck(), лучше kgo.AllISRAcks().
Один client на весь процесс - это правильно. Создавать новый client на каждый запрос - дорого (TCP handshake, получение метаданных). kgo.Client потокобезопасен.
Смотри также
- Message Brokers - Consumer Group - как читать то, что отправил producer
- Message Brokers - Exactly-once - idempotent producer и транзакции
- Event-Driven - Гарантии доставки - теория at-least-once и exactly-once
Мини-задание
- Подними Kafka из предыдущего урока, напиши
main.goсKafkaPublisherи отправь 5 событий - Проверь через
kafka-console-consumer --from-beginningчто все 5 дошли - Включи
ProducerLinger(5ms)и отправь 100 событий в цикле - сравни время с/без линтера - Попробуй отправить запись размером 2MB (больше
ProducerBatchMaxBytes) - что вернётProduceSync? - Реализуй
Close(ctx)сFlushи добавь его вhttp.Server.RegisterOnShutdown