karawaci.kode

← Semua snippet

Go Lanjut Otomasi

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