Implementing the Saga Pattern in Go Microservices: Orchestration and Choreography with PostgreSQL and Redis
Introduction: Why Distributed Transactions Are Hard
In a monolithic application, a transaction is an atomic operation within a single database. PostgreSQL's ACID guarantees make it reliable and predictable. But when a system is broken into microservices, each with its own database, classic transactions become impossible: you cannot issue a BEGIN ... COMMIT simultaneously across two isolated data stores.
Two-phase commit (2PC) formally addresses this problem, but in practice it introduces tight service coupling, blocking waits, and a single point of failure in the coordinator. For high-load systems, this is unacceptable.
The Saga pattern is an alternative approach: instead of one distributed transaction, a sequence of local transactions is executed in each service. If a failure occurs at any step, compensating transactions are triggered to undo the already-completed steps. This is not a rollback in the traditional sense — it consists of explicit reverse actions implemented in code.
Two Approaches: Orchestration vs. Choreography
The Saga pattern can be implemented in two fundamentally different ways. The choice between them shapes the architecture of the entire system.
Choreography
Each service publishes events after completing its local transaction. Other services subscribe to those events and react accordingly. There is no central coordinator — each service knows what to do in response to a specific event.
- Pros: full decentralization, low coupling, easy to add new participants.
- Cons: difficult to track the overall state of the saga, logic is spread across services, debugging becomes a detective exercise.
Orchestration
A central Saga Orchestrator manages all steps: it calls services, waits for responses, and decides when to trigger compensation. It stores the saga's state — typically in a database.
- Pros: full visibility into saga state, centralized logic, clear step ordering, straightforward debugging.
- Cons: the orchestrator is a potential single point of failure, tighter coupling to the coordinator, risk of the orchestrator becoming a "fat" service.
Comparing the Approaches
- Saga state: with choreography — distributed across events and services; with orchestration — centralized in the orchestrator's database.
- Coupling: choreography yields loose coupling between services; orchestration creates a dependency of all services on the orchestrator.
- Observability: choreography requires aggregating logs and traces; orchestration provides complete state from a single place.
- When to choose choreography: simple sagas (2–3 steps), the team already works with event-driven architecture, services are loosely coupled.
- When to choose orchestration: complex business processes (5+ steps), explicit state monitoring is required, centralized compensation is important.
Implementing a Choreography Saga in Go with Redis Streams
Redis Streams are a great fit for choreography — they provide a persistent message queue with consumer groups. Unlike simple Pub/Sub, messages are stored and can be re-read after a failure.
Event Structure
// events.go
package saga
import "time"
type EventType string
const (
OrderCreated EventType = "order.created"
PaymentProcessed EventType = "payment.processed"
PaymentFailed EventType = "payment.failed"
InventoryReserved EventType = "inventory.reserved"
InventoryFailed EventType = "inventory.failed"
OrderCompleted EventType = "order.completed"
OrderCancelled EventType = "order.cancelled"
)
type SagaEvent struct {
SagaID string `json:"saga_id"`
EventType EventType `json:"event_type"`
Payload map[string]string `json:"payload"`
Timestamp time.Time `json:"timestamp"`
}
Publishing Events to Redis Streams
// publisher.go
package saga
import (
"context"
"encoding/json"
"fmt"
"time"
"github.com/redis/go-redis/v9"
)
type EventPublisher struct {
rdb *redis.Client
stream string
}
func NewEventPublisher(rdb *redis.Client, stream string) *EventPublisher {
return &EventPublisher{rdb: rdb, stream: stream}
}
func (p *EventPublisher) Publish(ctx context.Context, event SagaEvent) error {
event.Timestamp = time.Now()
data, err := json.Marshal(event)
if err != nil {
return fmt.Errorf("marshal event: %w", err)
}
return p.rdb.XAdd(ctx, &redis.XAddArgs{
Stream: p.stream,
Values: map[string]interface{}{
"saga_id": event.SagaID,
"event_type": string(event.EventType),
"data": string(data),
},
}).Err()
}
Event Consumer with Consumer Group
// consumer.go
package saga
import (
"context"
"encoding/json"
"fmt"
"log"
"time"
"github.com/redis/go-redis/v9"
)
type EventHandler func(ctx context.Context, event SagaEvent) error
type EventConsumer struct {
rdb *redis.Client
stream string
group string
consumer string
handlers map[EventType]EventHandler
}
func NewEventConsumer(rdb *redis.Client, stream, group, consumer string) *EventConsumer {
return &EventConsumer{
rdb: rdb,
stream: stream,
group: group,
consumer: consumer,
handlers: make(map[EventType]EventHandler),
}
}
func (c *EventConsumer) Register(eventType EventType, handler EventHandler) {
c.handlers[eventType] = handler
}
func (c *EventConsumer) Start(ctx context.Context) error {
// Create the consumer group if it doesn't exist
c.rdb.XGroupCreateMkStream(ctx, c.stream, c.group, "0")
for {
select {
case <-ctx.Done():
return ctx.Err()
default:
}
streams, err := c.rdb.XReadGroup(ctx, &redis.XReadGroupArgs{
Group: c.group,
Consumer: c.consumer,
Streams: []string{c.stream, ">"},
Count: 10,
Block: 5 * time.Second,
}).Result()
if err != nil && err != redis.Nil {
log.Printf("xreadgroup error: %v", err)
continue
}
for _, stream := range streams {
for _, msg := range stream.Messages {
if err := c.processMessage(ctx, msg); err != nil {
log.Printf("process message %s error: %v", msg.ID, err)
continue // message remains in PEL for retry
}
// ACK only on successful processing
c.rdb.XAck(ctx, c.stream, c.group, msg.ID)
}
}
}
}
func (c *EventConsumer) processMessage(ctx context.Context, msg redis.XMessage) error {
dataStr, ok := msg.Values["data"].(string)
if !ok {
return fmt.Errorf("invalid message format")
}
var event SagaEvent
if err := json.Unmarshal([]byte(dataStr), &event); err != nil {
return fmt.Errorf("unmarshal event: %w", err)
}
handler, ok := c.handlers[event.EventType]
if !ok {
return nil // ignore unknown events
}
return handler(ctx, event)
}
Example: Payment Service in a Choreography Saga
// payment_service.go
package payment
import (
"context"
"fmt"
"log"
"github.com/yourorg/saga"
)
type PaymentService struct {
publisher *saga.EventPublisher
repo *PaymentRepository
}
func (s *PaymentService) HandleOrderCreated(ctx context.Context, event saga.SagaEvent) error {
orderID := event.Payload["order_id"]
amount := event.Payload["amount"]
// Idempotency check — do not process twice
if exists, _ := s.repo.ExistsPayment(ctx, event.SagaID); exists {
log.Printf("payment for saga %s already processed, skipping", event.SagaID)
return nil
}
err := s.repo.ProcessPayment(ctx, event.SagaID, orderID, amount)
if err != nil {
// Publish a failure event — other services will start compensation
return s.publisher.Publish(ctx, saga.SagaEvent{
SagaID: event.SagaID,
EventType: saga.PaymentFailed,
Payload: map[string]string{"reason": err.Error()},
})
}
return s.publisher.Publish(ctx, saga.SagaEvent{
SagaID: event.SagaID,
EventType: saga.PaymentProcessed,
Payload: map[string]string{"order_id": orderID},
})
}
func (s *PaymentService) HandleOrderCancelled(ctx context.Context, event saga.SagaEvent) error {
// Compensating transaction: refund the payment
if err := s.repo.RefundPayment(ctx, event.SagaID); err != nil {
return fmt.Errorf("refund payment: %w", err)
}
log.Printf("payment refunded for saga %s", event.SagaID)
return nil
}
Implementing an Orchestration Saga in Go with PostgreSQL
The orchestrator stores the state of each saga in PostgreSQL. This enables recovery after a failure: on restart, the orchestrator reads incomplete sagas and resumes execution.
Database Schema
-- migrations/001_saga_state.sql
CREATE TYPE saga_status AS ENUM (
'started',
'processing',
'completed',
'compensating',
'failed'
);
CREATE TABLE saga_instances (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
saga_type VARCHAR(100) NOT NULL,
status saga_status NOT NULL DEFAULT 'started',
current_step VARCHAR(100),
payload JSONB NOT NULL DEFAULT '{}',
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
);
CREATE TABLE saga_steps (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
saga_id UUID NOT NULL REFERENCES saga_instances(id),
step_name VARCHAR(100) NOT NULL,
status VARCHAR(50) NOT NULL DEFAULT 'pending',
result JSONB,
executed_at TIMESTAMPTZ,
UNIQUE(saga_id, step_name)
);
CREATE INDEX idx_saga_instances_status ON saga_instances(status);
State Model and Repository
// orchestrator/state.go
package orchestrator
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"time"
"github.com/google/uuid"
)
type SagaStatus string
const (
StatusStarted SagaStatus = "started"
StatusProcessing SagaStatus = "processing"
StatusCompleted SagaStatus = "completed"
StatusCompensating SagaStatus = "compensating"
StatusFailed SagaStatus = "failed"
)
type SagaInstance struct {
ID uuid.UUID `db:"id"`
SagaType string `db:"saga_type"`
Status SagaStatus `db:"status"`
CurrentStep string `db:"current_step"`
Payload map[string]any `db:"payload"`
CreatedAt time.Time `db:"created_at"`
UpdatedAt time.Time `db:"updated_at"`
}
type SagaRepository struct {
db *sql.DB
}
func (r *SagaRepository) Create(ctx context.Context, sagaType string, payload map[string]any) (*SagaInstance, error) {
payloadJSON, err := json.Marshal(payload)
if err != nil {
return nil, err
}
instance := &SagaInstance{
ID: uuid.New(),
SagaType: sagaType,
Status: StatusStarted,
Payload: payload,
}
_, err = r.db.ExecContext(ctx,
`INSERT INTO saga_instances (id, saga_type, status, payload)
VALUES ($1, $2, $3, $4)`,
instance.ID, instance.SagaType, instance.Status, payloadJSON,
)
return instance, err
}
func (r *SagaRepository) UpdateStatus(ctx context.Context, sagaID uuid.UUID, status SagaStatus, step string) error {
_, err := r.db.ExecContext(ctx,
`UPDATE saga_instances
SET status = $1, current_step = $2, updated_at = NOW()
WHERE id = $3`,
status, step, sagaID,
)
return err
}
func (r *SagaRepository) MarkStepDone(ctx context.Context, sagaID uuid.UUID, stepName string, result map[string]any) error {
resultJSON, _ := json.Marshal(result)
_, err := r.db.ExecContext(ctx,
`INSERT INTO saga_steps (saga_id, step_name, status, result, executed_at)
VALUES ($1, $2, 'completed', $3, NOW())
ON CONFLICT (saga_id, step_name) DO UPDATE
SET status = 'completed', result = $3, executed_at = NOW()`,
sagaID, stepName, resultJSON,
)
return err
}
func (r *SagaRepository) IsStepDone(ctx context.Context, sagaID uuid.UUID, stepName string) (bool, error) {
var count int
err := r.db.QueryRowContext(ctx,
`SELECT COUNT(*) FROM saga_steps
WHERE saga_id = $1 AND step_name = $2 AND status = 'completed'`,
sagaID, stepName,
).Scan(&count)
return count > 0, err
}
func (r *SagaRepository) GetPendingSagas(ctx context.Context) ([]*SagaInstance, error) {
rows, err := r.db.QueryContext(ctx,
`SELECT id, saga_type, status, current_step, payload, created_at, updated_at
FROM saga_instances
WHERE status IN ('started', 'processing', 'compensating')
AND updated_at < NOW() - INTERVAL '5 minutes'`)
if err != nil {
return nil, fmt.Errorf("query pending sagas: %w", err)
}
defer rows.Close()
var instances []*SagaInstance
for rows.Next() {
var inst SagaInstance
var payloadJSON []byte
if err := rows.Scan(&inst.ID, &inst.SagaType, &inst.Status,
&inst.CurrentStep, &payloadJSON, &inst.CreatedAt, &inst.UpdatedAt); err != nil {
return nil, err
}
json.Unmarshal(payloadJSON, &inst.Payload)
instances = append(instances, &inst)
}
return instances, rows.Err()
}
Central Orchestrator
// orchestrator/orchestrator.go
package orchestrator
import (
"context"
"fmt"
"log"
"github.com/google/uuid"
)
type Step struct {
Name string
Execute func(ctx context.Context, payload map[string]any) (map[string]any, error)
Compensate func(ctx context.Context, payload map[string]any) error
}
type SagaOrchestrator struct {
repo *SagaRepository
steps []Step
}
func NewSagaOrchestrator(repo *SagaRepository, steps []Step) *SagaOrchestrator {
return &SagaOrchestrator{repo: repo, steps: steps}
}
func (o *SagaOrchestrator) Execute(ctx context.Context, sagaType string, payload map[string]any) error {
instance, err := o.repo.Create(ctx, sagaType, payload)
if err != nil {
return fmt.Errorf("create saga: %w", err)
}
return o.run(ctx, instance)
}
func (o *SagaOrchestrator) run(ctx context.Context, instance *SagaInstance) error {
executedSteps := []int{}
for i, step := range o.steps {
// Idempotency check: already completed?
done, err := o.repo.IsStepDone(ctx, instance.ID, step.Name)
if err != nil {
return fmt.Errorf("check step %s: %w", step.Name, err)
}
if done {
executedSteps = append(executedSteps, i)
continue
}
o.repo.UpdateStatus(ctx, instance.ID, StatusProcessing, step.Name)
log.Printf("[saga %s] executing step: %s", instance.ID, step.Name)
result, err := step.Execute(ctx, instance.Payload)
if err != nil {
log.Printf("[saga %s] step %s failed: %v — starting compensation", instance.ID, step.Name, err)
o.repo.UpdateStatus(ctx, instance.ID, StatusCompensating, step.Name)
return o.compensate(ctx, instance, executedSteps)
}
// Save the step result for idempotency
if err := o.repo.MarkStepDone(ctx, instance.ID, step.Name, result); err != nil {
return fmt.Errorf("mark step done: %w", err)
}
executedSteps = append(executedSteps, i)
// Update the payload with the step's results
for k, v := range result {
instance.Payload[k] = v
}
}
return o.repo.UpdateStatus(ctx, instance.ID, StatusCompleted, "")
}
func (o *SagaOrchestrator) compensate(ctx context.Context, instance *SagaInstance, executedSteps []int) error {
// Compensate in reverse order
for i := len(executedSteps) - 1; i >= 0; i-- {
stepIdx := executedSteps[i]
step := o.steps[stepIdx]
if step.Compensate == nil {
continue
}
log.Printf("[saga %s] compensating step: %s", instance.ID, step.Name)
if err := step.Compensate(ctx, instance.Payload); err != nil {
// Log the error but continue compensating other steps
log.Printf("[saga %s] compensation of %s failed: %v", instance.ID, step.Name, err)
}
}
return o.repo.UpdateStatus(ctx, instance.ID, StatusFailed, "")
}
Using the Orchestrator
// main.go (assembly example)
package main
import (
"context"
"database/sql"
"log"
_ "github.com/lib/pq"
"github.com/yourorg/orchestrator"
)
func main() {
db, _ := sql.Open("postgres", "postgres://user:pass@localhost/orders?sslmode=disable")
repo := &orchestrator.SagaRepository{} // initialize with db
orderSaga := orchestrator.NewSagaOrchestrator(repo, []orchestrator.Step{
{
Name: "reserve_inventory",
Execute: func(ctx context.Context, payload map[string]any) (map[string]any, error) {
// call inventory service
return map[string]any{"reservation_id": "res-123"}, nil
},
Compensate: func(ctx context.Context, payload map[string]any) error {
// cancel reservation
return nil
},
},
{
Name: "process_payment",
Execute: func(ctx context.Context, payload map[string]any) (map[string]any, error) {
// call payment service
return map[string]any{"payment_id": "pay-456"}, nil
},
Compensate: func(ctx context.Context, payload map[string]any) error {
// refund payment
return nil
},
},
{
Name: "confirm_order",
Execute: func(ctx context.Context, payload map[string]any) (map[string]any, error) {
return map[string]any{"status": "confirmed"}, nil
},
Compensate: nil, // final step, no compensation needed
},
})
ctx := context.Background()
if err := orderSaga.Execute(ctx, "order_saga", map[string]any{
"order_id": "ord-789",
"amount": "150.00",
"user_id": "usr-001",
}); err != nil {
log.Printf("saga failed: %v", err)
}
}
Handling Compensating Transactions and Idempotency
Compensating transactions are not rollbacks. They are separate business operations that must be implemented explicitly. Several critically important rules apply:
Idempotency Is a Hard Requirement
Any operation in a saga — both forward and compensating — may be called more than once due to retries after network failures. Every handler must behave correctly when invoked multiple times.
The standard pattern is to store an idempotency_key (typically saga_id + step_name) in a results table. Before executing the operation, check whether a record with that key already exists.
// idempotency.go
func (r *PaymentRepository) ProcessPaymentIdempotent(
ctx context.Context,
idempotencyKey string,
amount string,
) error {
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
return err
}
defer tx.Rollback()
// Atomic check + insert via INSERT ... ON CONFLICT DO NOTHING
result, err := tx.ExecContext(ctx,
`INSERT INTO payment_operations (idempotency_key, amount, created_at)
VALUES ($1, $2, NOW())
ON CONFLICT (idempotency_key) DO NOTHING`,
idempotencyKey, amount,
)
if err != nil {
return fmt.Errorf("insert payment operation: %w", err)
}
rows, _ := result.RowsAffected()
if rows == 0 {
// Operation already completed — return success
return tx.Commit()
}
// Execute the actual payment
if err := r.chargeCard(ctx, tx, amount); err != nil {
return err
}
return tx.Commit()
}
Retries with Exponential Backoff
For transient failures (network issues, database overload), use retries with jitter to avoid thundering herds:
// retry.go
package saga
import (
"context"
"math"
"math/rand"
"time"
)
func WithRetry(ctx context.Context, maxAttempts int, fn func() error) error {
var lastErr error
for attempt := 0; attempt < maxAttempts; attempt++ {
if err := fn(); err == nil {
return nil
} else {
lastErr = err
}
if attempt == maxAttempts-1 {
break
}
backoff := time.Duration(math.Pow(2, float64(attempt))) * 100 * time.Millisecond
jitter := time.Duration(rand.Intn(100)) * time.Millisecond
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(backoff + jitter):
}
}
return lastErr
}
Monitoring and Debugging Saga Transactions
Without observability, sagas become a black box. For production systems, you need:
- Structured logging: every saga event is logged with fields
saga_id,step,status,duration_ms. Uselog/slogorzerolog. - Metrics: counters for started, completed, and compensated sagas via Prometheus. Separate histograms for the duration of each step.
- Distributed tracing: propagate
trace_idacross all services through the context. OpenTelemetry + Jaeger provide a complete picture of how a saga flows through services. - Dead Letter Queue (DLQ): in Redis Streams, messages that have failed repeatedly should be moved to a separate stream for manual analysis.
- Admin endpoint: a simple HTTP endpoint that returns the saga state by its ID from PostgreSQL — invaluable during incidents.
// monitoring.go
package orchestrator
import (
"log/slog"
"time"
"github.com/prometheus/client_golang/prometheus"
)
var (
sagaTotal = prometheus.NewCounterVec(
prometheus.CounterOpts{Name: "saga_total", Help: "Total saga executions"},
[]string{"type", "status"},
)
sagaDuration = prometheus.NewHistogramVec(
prometheus.HistogramOpts{
Name: "saga_duration_seconds",
Help: "Saga execution duration",
Buckets: prometheus.DefBuckets,
},
[]string{"type"},
)
)
func init() {
prometheus.MustRegister(sagaTotal, sagaDuration)
}
func recordSagaMetrics(sagaType string, status SagaStatus, start time.Time) {
sagaTotal.WithLabelValues(sagaType, string(status)).Inc()
sagaDuration.WithLabelValues(sagaType).Observe(time.Since(start).Seconds())
slog.Info("saga finished",
"type", sagaType,
"status", status,
"duration_ms", time.Since(start).Milliseconds(),
)
}
Common Mistakes and Anti-Patterns
- Non-idempotent handlers. The most frequent mistake. If a saga step is not idempotent, a repeated call will create duplicate operations — double charges, double reservations. Always use an
idempotency_key. - Missing timeouts. A saga step can hang indefinitely waiting for a response from an unavailable service. Every call must have a context with a deadline. Use
context.WithTimeout. - Compensation without retries. If a compensating transaction fails, that is a critical situation. Retries, alerts, and possibly manual intervention are required. Never silently ignore compensation errors.
- Overly long sagas. A saga with 10+ steps signals problems in service decomposition. A long saga means a prolonged window of data inconsistency and a massive compensation graph.
- Reading uncommitted data. While a saga is executing, data is in an intermediate state. Other services reading that data may get an incorrect picture. Design for eventual consistency explicitly.
- Storing saga state only in memory. After a service restart, all in-flight sagas are lost. State must be persisted to PostgreSQL before each step begins.
- Ignoring compensation order. Compensation must happen in the strict reverse order of the executed steps. Otherwise, new inconsistencies can be introduced.
Conclusion: Choosing the Right Approach
The Saga pattern addresses a real problem of distributed transactions in Go microservice systems, but it requires careful design and disciplined implementation.
Saga doesn't make the system simpler — it makes complexity explicit and manageable.
Choose choreography if you have 2–4 participants in the saga, your team already works with an event-driven approach and Redis Streams, and you are willing to invest in good tracing for debugging. Choreography scales horizontally with ease.
Choose orchestration if the saga involves complex business logic with branching, you need full state visibility in PostgreSQL, the team prefers explicit flow control, and fast incident debugging is critical.
Regardless of your choice: implement idempotency from day one, add distributed tracing before going to production, provide a mechanism for recovering stalled sagas through periodic database scanning, and never ignore errors from compensating transactions. Go, with its explicit error handling and powerful concurrency primitives, is an excellent choice for implementing either approach.
Technologies
Tags
Ruslan Ismailov
Senior Web / Backend Developer. Senior web/backend developer with 9 years of experience. Stack: PHP, Laravel, PostgreSQL, Redis, Docker, Kubernetes, REST, microservices, CI/CD. More about me →