Idempotent Kafka Consumer (Go)
GoSkip events that were already processed, inside one transaction.
Go · 19 lines
copied 10 times
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