Consumer Group на Go: партиции, offset commit

Consumer Group на Go: партиции, offset commit

Consumer Group - механизм горизонтального масштабирования в Kafka. Несколько воркеров объединяются под одним group.id, Kafka распределяет партиции между ними. Добавляешь воркер - автоматически забирает часть партиций. Убираешь - его партиции переходят к оставшимся (rebalance).

Базовый consumer loop

package main

import (
    "context"
    "log"
    "os/signal"
    "syscall"

    "github.com/twmb/franz-go/pkg/kgo"
)

func main() {
    ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
    defer stop()

    client, err := kgo.NewClient(
        kgo.SeedBrokers("localhost:9092"),
        kgo.ConsumerGroup("analytics"),
        kgo.ConsumeTopics("lesson-events"),
    )
    if err != nil {
        log.Fatal(err)
    }
    defer client.Close()

    for {
        fetches := client.PollFetches(ctx)
        if fetches.IsClientClosed() || ctx.Err() != nil {
            return
        }
        if errs := fetches.Errors(); len(errs) > 0 {
            for _, e := range errs {
                log.Printf("fetch error: topic=%s partition=%d err=%v", e.Topic, e.Partition, e.Err)
            }
            continue
        }

        fetches.EachRecord(func(record *kgo.Record) {
            handleEvent(record)
        })
    }
}

func handleEvent(r *kgo.Record) {
    log.Printf("partition=%d offset=%d key=%s value=%s",
        r.Partition, r.Offset, r.Key, r.Value)
    // бизнес-логика
}

По умолчанию franz-go делает auto-commit каждые 5 секунд. При краше после обработки и до commit - события будут переданы повторно (at-least-once).

Manual Commit для надёжной обработки

client, err := kgo.NewClient(
    kgo.SeedBrokers("localhost:9092"),
    kgo.ConsumerGroup("analytics"),
    kgo.ConsumeTopics("lesson-events"),
    kgo.DisableAutoCommit(),
)

for {
    fetches := client.PollFetches(ctx)
    if ctx.Err() != nil {
        return
    }

    var toCommit []*kgo.Record

    fetches.EachRecord(func(r *kgo.Record) {
        if err := process(r); err != nil {
            log.Printf("process error: %v - skipping", err)
            return
        }
        toCommit = append(toCommit, r)
    })

    if len(toCommit) > 0 {
        if err := client.CommitRecords(ctx, toCommit...); err != nil {
            log.Printf("commit error: %v", err)
        }
    }
}

CommitRecords коммитит offset для каждой партиции по максимальному offset в переданном списке.

Авто-коммит: что за ним скрывается

Auto-commit коммитит последний опрошенный offset, а не последний обработанный. Это значит:

PollFetches → получены записи 100-120
  обработка записей 100-115...
  краш приложения
  auto-commit НЕ успел сработать (или сработал до краша)

→ при перезапуске: читаем снова с 100 (ok, at-least-once)
  ИЛИ читаем с 121 (плохо! 116-120 потеряны)

Для бизнес-критичной обработки - всегда manual commit после успешной обработки.

Параллельная обработка

fetches := client.PollFetches(ctx)

var wg sync.WaitGroup
var mu sync.Mutex
var toCommit []*kgo.Record

fetches.EachPartition(func(p kgo.FetchTopicPartition) {
    // Обрабатываем каждую партицию в своей горутине
    // Порядок внутри партиции сохраняется
    wg.Add(1)
    go func(p kgo.FetchTopicPartition) {
        defer wg.Done()
        p.EachRecord(func(r *kgo.Record) {
            if err := process(r); err == nil {
                mu.Lock()
                toCommit = append(toCommit, r)
                mu.Unlock()
            }
        })
    }(p)
})

wg.Wait()

if len(toCommit) > 0 {
    client.CommitRecords(ctx, toCommit...)
}

Разные партиции обрабатываются параллельно (порядок между ними всё равно не гарантирован). Внутри одной партиции - последовательно.

Rebalance: hooks

Kafka выполняет rebalance при добавлении/удалении consumer-а. На время rebalance все consumer-ы в группе останавливают чтение. franz-go уведомляет через хуки:

client, _ := kgo.NewClient(
    kgo.SeedBrokers("localhost:9092"),
    kgo.ConsumerGroup("analytics"),
    kgo.ConsumeTopics("lesson-events"),
    kgo.OnPartitionsAssigned(func(ctx context.Context, c *kgo.Client, assigned map[string][]int32) {
        for topic, partitions := range assigned {
            log.Printf("assigned: topic=%s partitions=%v", topic, partitions)
        }
    }),
    kgo.OnPartitionsRevoked(func(ctx context.Context, c *kgo.Client, revoked map[string][]int32) {
        // Обязательно commit перед тем как отдать партицию
        c.CommitMarkedOffsets(ctx)
        log.Printf("revoked partitions, committed offsets")
    }),
)

OnPartitionsRevoked - правильное место для финального commit. Без него при rebalance возможна повторная обработка.

Consumer Lag

Consumer lag = latest_offset - committed_offset. Растущий lag означает, что consumer не успевает за producer-ом.

# Проверить lag через CLI
docker exec kafka kafka-consumer-groups \
  --bootstrap-server localhost:9092 \
  --describe \
  --group analytics

# Вывод:
# GROUP     TOPIC          PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
# analytics lesson-events  0          1000            1050            50
# analytics lesson-events  1          800             800             0
# analytics lesson-events  2          600             650             50

Lag > 0 - нормально при кратких всплесках. Постоянно растущий lag - нужно масштабировать consumer group (добавить воркеры) или оптимизировать обработку.

Graceful Shutdown

ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()

go func() {
    <-ctx.Done()
    log.Println("shutdown signal received")
    // PollFetches вернёт после отмены контекста
}()

for {
    fetches := client.PollFetches(ctx)
    if ctx.Err() != nil {
        // Commit последних обработанных записей перед выходом
        if err := client.CommitMarkedOffsets(context.Background()); err != nil {
            log.Printf("final commit error: %v", err)
        }
        return
    }
    // обработка...
}

client.Close() отправляет heartbeat с LeaveGroup запросом - Kafka выполняет rebalance немедленно, а не ждёт сессионный таймаут (по умолчанию 45 секунд).

Consume без group (standalone)

// Читать с начала, без consumer group - для утилит и отладки
client, _ := kgo.NewClient(
    kgo.SeedBrokers("localhost:9092"),
    kgo.ConsumeTopics("lesson-events"),
    kgo.ConsumeResetOffset(kgo.NewOffset().AtStart()),
)

// Или читать конкретный диапазон партиций
client, _ := kgo.NewClient(
    kgo.SeedBrokers("localhost:9092"),
    kgo.ConsumePartitions(map[string]map[int32]kgo.Offset{
        "lesson-events": {
            0: kgo.NewOffset().At(100), // с offset 100 партиции 0
        },
    }),
)

Без ConsumerGroup - нет координации, нет rebalance. Один процесс читает все партиции сам.

Типичные ошибки

Больше воркеров чем партиций - лишние воркеры простаивают и не получают партиций. При 3 партициях максимум 3 активных consumer-а в группе.

Не обрабатывать ctx.Err() после PollFetches - при отмене контекста PollFetches возвращает пустой Fetches. Без проверки ctx.Err() цикл крутится бесконечно.

Commit до обработки - at-most-once. Если коммитить offset сразу после poll, а потом обрабатывать - при краше во время обработки сообщение потеряно.

Глобальный time.Sleep в обработчике - блокирует партицию полностью. Используй отдельные воркеры per-partition или буферизованный канал.

Смотри также

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

  • Напиши consumer с DisableAutoCommit() и manual CommitRecords после обработки каждого батча
  • Запусти двух воркеров с одним group.id на топике с 3 партициями - выведи в логах какая партиция у каждого
  • Остановь один воркер - через сколько секунд партиции перераспределились ко второму?
  • Отправь 100 событий через producer, затем остановь consumer после 50 - убедись что при перезапуске consumer прочитал ровно оставшиеся 50
  • Добавь OnPartitionsRevoked hook с CommitMarkedOffsets

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