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.startedlesson.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 потокобезопасен.

Смотри также

Мини-задание

  • Подними Kafka из предыдущего урока, напиши main.go с KafkaPublisher и отправь 5 событий
  • Проверь через kafka-console-consumer --from-beginning что все 5 дошли
  • Включи ProducerLinger(5ms) и отправь 100 событий в цикле - сравни время с/без линтера
  • Попробуй отправить запись размером 2MB (больше ProducerBatchMaxBytes) - что вернёт ProduceSync?
  • Реализуй Close(ctx) с Flush и добавь его в http.Server.RegisterOnShutdown

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.