Skip to main content

Idempotent Kafka Consumer (Go)

Go

Skip events that were already processed, inside one transaction.

L
Lindiwe Dube
Sample Oct 2, 2026
Go · 19 lines
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
18
19
func handle(ctx context.Context, db *sql.DB, evt Event) error {
	tx, err := db.BeginTx(ctx, nil)
	if err != nil {
		return err
	}
	defer tx.Rollback()
	res, err := tx.ExecContext(ctx,
		`INSERT INTO processed_events (id) VALUES ($1) ON CONFLICT DO NOTHING`, evt.ID)
	if err != nil {
		return err
	}
	if n, _ := res.RowsAffected(); n == 0 {
		return nil // already handled
	}
	if err := apply(ctx, tx, evt); err != nil {
		return err
	}
	return tx.Commit()
}
2 Fire 1 Learned something 0 Saved me time 3 Mind-blown
Sign in to react and star