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 или буферизованный канал.
Смотри также
- Message Brokers - Producer - как отправлять то, что читаем
- Message Brokers - Exactly-once - commit стратегия и гарантии
- Go - Context - отмена через context, signal.NotifyContext
Мини-задание
- Напиши consumer с
DisableAutoCommit()и manualCommitRecordsпосле обработки каждого батча - Запусти двух воркеров с одним
group.idна топике с 3 партициями - выведи в логах какая партиция у каждого - Остановь один воркер - через сколько секунд партиции перераспределились ко второму?
- Отправь 100 событий через producer, затем остановь consumer после 50 - убедись что при перезапуске consumer прочитал ровно оставшиеся 50
- Добавь
OnPartitionsRevokedhook сCommitMarkedOffsets