Go Kafka consumer group dengan rebalance
Consumer group Kafka di Go pakai franz-go — handle rebalance, graceful shutdown, manual commit pattern. Production-ready event consumer.
Dipublikasikan 17 Juli 2026
Kafka consumer di Go yang tidak handle rebalance properly = duplicate message saat scaling atau deploy. franz-go (kgo) lebih modern dari sarama tapi pattern-nya berbeda. Snippet ini consumer group untuk event “order.created” Tokopedia dengan manual commit + graceful drain.
Kode
// Cargo equivalent untuk Go — go.mod
// require (
// github.com/twmb/franz-go v1.18.0
// github.com/twmb/franz-go/pkg/kadm v1.13.0
// )
package main
import (
"context"
"encoding/json"
"errors"
"fmt"
"log/slog"
"os"
"os/signal"
"sync"
"syscall"
"time"
"github.com/twmb/franz-go/pkg/kgo"
"github.com/twmb/franz-go/pkg/kmsg"
)
type OrderEvent struct {
OrderID int64 `json:"order_id"`
UserID int64 `json:"user_id"`
TotalRupiah int64 `json:"total_rupiah"`
Status string `json:"status"`
CreatedAt time.Time `json:"created_at"`
}
type Consumer struct {
client *kgo.Client
logger *slog.Logger
// Processed map untuk idempotency — production pakai Redis/DB
processed sync.Map
}
func NewConsumer(brokers []string, groupID, topic string, logger *slog.Logger) (*Consumer, error) {
client, err := kgo.NewClient(
kgo.SeedBrokers(brokers...),
kgo.ConsumerGroup(groupID),
kgo.ConsumeTopics(topic),
// Manual commit — kontrol penuh kapan offset di-advance
kgo.DisableAutoCommit(),
// Rebalance behavior
kgo.OnPartitionsAssigned(func(ctx context.Context, cl *kgo.Client, partitions map[string][]int32) {
logger.Info("partitions assigned", "partitions", partitions)
}),
kgo.OnPartitionsRevoked(func(ctx context.Context, cl *kgo.Client, partitions map[string][]int32) {
logger.Info("partitions revoked, commit final offset", "partitions", partitions)
// Important: commit semua processed sebelum partisi diambil consumer lain
if err := cl.CommitMarkedOffsets(ctx); err != nil {
logger.Error("commit on revoke fail", "err", err)
}
}),
kgo.OnPartitionsLost(func(ctx context.Context, cl *kgo.Client, partitions map[string][]int32) {
logger.Warn("partitions LOST (session timeout)", "partitions", partitions)
// Lost = sudah tidak punya hak commit. Drop in-flight work.
}),
// Reset offset di partition baru
kgo.ConsumeResetOffset(kgo.NewOffset().AtStart()),
// Session timeout — kalau gak heartbeat dalam 30s, di-evict
kgo.SessionTimeout(30*time.Second),
kgo.HeartbeatInterval(3*time.Second),
// Fetch behavior
kgo.FetchMaxBytes(50*1024*1024), // 50 MB max per fetch
kgo.FetchMaxWait(500*time.Millisecond),
)
if err != nil {
return nil, fmt.Errorf("new client: %w", err)
}
// Verify connection
if err := client.Ping(context.Background()); err != nil {
return nil, fmt.Errorf("ping: %w", err)
}
return &Consumer{client: client, logger: logger}, nil
}
func (c *Consumer) Run(ctx context.Context) error {
c.logger.Info("consumer started")
for {
// Poll batch — block max FetchMaxWait
fetches := c.client.PollFetches(ctx)
// Handle error per partition
if errs := fetches.Errors(); len(errs) > 0 {
for _, e := range errs {
if errors.Is(e.Err, context.Canceled) {
c.logger.Info("ctx canceled, drain & exit")
return nil
}
c.logger.Error("fetch error", "topic", e.Topic, "partition", e.Partition, "err", e.Err)
}
}
if fetches.Empty() {
continue
}
// Process per partition — preserve order dalam 1 partition
var wg sync.WaitGroup
fetches.EachPartition(func(p kgo.FetchTopicPartition) {
wg.Add(1)
go func(partition kgo.FetchTopicPartition) {
defer wg.Done()
c.processPartition(ctx, partition)
}(p)
})
wg.Wait()
// Commit semua yang sukses (yang sudah ditandai via MarkCommitRecords)
if err := c.client.CommitMarkedOffsets(ctx); err != nil {
c.logger.Error("commit fail", "err", err)
}
}
}
func (c *Consumer) processPartition(ctx context.Context, p kgo.FetchTopicPartition) {
for _, record := range p.Records {
select {
case <-ctx.Done():
return
default:
}
// Idempotency check via key — penting kalau ada duplicate
key := string(record.Key)
if _, seen := c.processed.LoadOrStore(key, struct{}{}); seen {
c.logger.Debug("skip duplicate", "key", key, "offset", record.Offset)
c.client.MarkCommitRecords(record)
continue
}
if err := c.processRecord(ctx, record); err != nil {
c.logger.Error("process fail",
"topic", record.Topic,
"partition", record.Partition,
"offset", record.Offset,
"err", err,
)
// Decision: kirim ke DLQ vs retry
if errors.Is(err, errPermanent) {
if err := c.sendToDLQ(ctx, record); err != nil {
c.logger.Error("DLQ send fail", "err", err)
// Tetap mark commit supaya tidak stuck
}
c.client.MarkCommitRecords(record)
}
// transient error — jangan mark, akan re-poll setelah rebalance / restart
return // stop process di partisi ini
}
c.client.MarkCommitRecords(record)
}
}
var errPermanent = errors.New("permanent error")
func (c *Consumer) processRecord(ctx context.Context, record *kgo.Record) error {
var event OrderEvent
if err := json.Unmarshal(record.Value, &event); err != nil {
return fmt.Errorf("%w: invalid JSON: %s", errPermanent, err)
}
c.logger.Info("processing order",
"order_id", event.OrderID,
"total", event.TotalRupiah,
"partition", record.Partition,
"offset", record.Offset,
)
// Business logic: kirim notifikasi, update analytics, etc
// HARUS idempotent — Kafka at-least-once delivery
if event.TotalRupiah <= 0 {
return fmt.Errorf("%w: total invalid", errPermanent)
}
// Simulasi work
time.Sleep(50 * time.Millisecond)
return nil
}
func (c *Consumer) sendToDLQ(ctx context.Context, record *kgo.Record) error {
dlqRecord := &kgo.Record{
Topic: record.Topic + ".dlq",
Key: record.Key,
Value: record.Value,
Headers: append(record.Headers,
kmsg.RecordHeader{Key: "x-original-topic", Value: []byte(record.Topic)},
kmsg.RecordHeader{Key: "x-original-offset", Value: []byte(fmt.Sprintf("%d", record.Offset))},
),
}
return c.client.ProduceSync(ctx, dlqRecord).FirstErr()
}
func (c *Consumer) Close() {
c.logger.Info("closing consumer (will trigger rebalance)")
c.client.Close()
}
func main() {
logger := slog.New(slog.NewJSONHandler(os.Stdout, nil))
consumer, err := NewConsumer(
[]string{"kafka-1:9092", "kafka-2:9092", "kafka-3:9092"},
"order-processor-v1",
"order.created",
logger,
)
if err != nil {
logger.Error("init consumer", "err", err)
os.Exit(1)
}
defer consumer.Close()
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGTERM, syscall.SIGINT)
defer stop()
if err := consumer.Run(ctx); err != nil {
logger.Error("consumer run", "err", err)
os.Exit(1)
}
logger.Info("shutdown clean")
}
Pemakaian
# Run consumer
go run main.go
# Output JSON log:
# {"time":"2026-07-17T10:00:00Z","level":"INFO","msg":"consumer started"}
# {"time":"2026-07-17T10:00:01Z","level":"INFO","msg":"partitions assigned","partitions":{"order.created":[0,1,2]}}
# {"time":"2026-07-17T10:00:02Z","level":"INFO","msg":"processing order","order_id":12345,"total":250000,"partition":0,"offset":154823}
# Scale up — deploy 1 instance lagi
docker run -d order-consumer:latest
# Log akan show rebalance otomatis:
# {"msg":"partitions revoked","partitions":{"order.created":[0,1,2]}}
# {"msg":"partitions assigned","partitions":{"order.created":[0,1]}} ← cuma punya 2 dari 3
# Cek consumer group dari kafka tools
kafka-consumer-groups.sh --bootstrap-server kafka-1:9092 \
--group order-processor-v1 --describe
# GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
# order-processor-v1 order.created 0 152031 152089 58
# order-processor-v1 order.created 1 148229 148245 16
# order-processor-v1 order.created 2 151102 151102 0
Kapan dipakai
- Microservice consume event dari Kafka.
- CDC pipeline (Debezium → consumer).
- Order processing system — fan-out ke service hilir.
- Log aggregation (consume log topic → push ke storage).
Catatan
- DisableAutoCommit + MarkCommitRecords — pattern manual commit. Kafka at-least-once delivery, tanpa idempotency ada duplicate.
- OnPartitionsRevoked wajib commit final offset. Tanpa ini, partition baru consumer mulai dari offset terakhir di-commit (mungkin lama) → duplicate.
- OnPartitionsLost — beda dengan revoked. Lost = session timeout, sudah tidak punya hak commit. Drop in-flight work.
- Idempotency key — pakai key Kafka atau field business (order_id). Tanpa idempotency, retry akan double-process.
- DLQ pattern — permanent error (invalid JSON, business rule fail) kirim ke DLQ topic. Transient error (DB down) jangan ke DLQ, biar retry naturally.
- Session timeout vs heartbeat — heartbeat 3s, session 30s. Network blip 10s tidak trigger rebalance.
- Fetch max bytes — terlalu besar = memory pressure, terlalu kecil = banyak round trip. 50MB OK untuk most case.
Jangan process record dalam goroutine terpisah dari poll loop — order dalam partition jadi tidak terpreserve. Pattern di snippet ini per-partition goroutine, cross-partition paralel.
# tags
gokafkaconsumer-groupfranz-gomessaging
Ditulis oleh Asti Larasati · 17 Juli 2026