Architecting Distributed Sagas with Temporal and Go
Architecting Distributed Sagas with Temporal and Go
Introduction
In a distributed microservices architecture, maintaining transactional consistency across multiple service boundaries is a notorious engineering challenge. Traditional database-level transactions (2PC - Two-Phase Commit) do not scale in highly distributed cloud-native environments due to network latency, locking overhead, and tight coupling between systems.
To solve this, architects use the Saga Pattern—a sequence of local transactions where each step updates a service-specific database. If a step fails, the Saga runs a series of compensating transactions in reverse order to undo the changes. While choreography-based Sagas (using message brokers like Kafka) are common, they often lead to "spaghetti event flows" that are difficult to trace and debug. This article demonstrates how to architect an orchestrator-based Saga using Temporal and Go, providing deterministic, auditable, and resilient distributed transactions.
Core Concepts & Architecture
Temporal simplifies distributed systems by letting developers write stateful, long-running workflows in native code. A Temporal-based Saga utilizes three primary abstractions:
-
Temporal Workflow: A deterministic orchestrator function that defines the sequence of execution steps and registers compensations.
-
Temporal Activities: Non-deterministic, side-effect-heavy tasks (such as calling external REST APIs, charging credit cards, or querying databases).
-
Saga Compensations: Registered rollback activities executed sequentially if any phase of the workflow fails.
System Topology
[ Client ] -> [ Temporal Client (Go API Gateway) ] │ ▼ [ Temporal Cluster ] (State Storage, Timers, Queues) │ ┌──────────┴──────────┐ ▼ ▼ [ Worker: Order Svc ] [ Worker: Payment Svc ]
-
CreateOrderActivity - AuthorizePaymentActivity
-
CancelOrderActivity - RefundPaymentActivity
In this model, the Temporal Cluster acts as the durable state machine. It schedules task executions, manages retries, and maintains the audit log, while the Go Workers pull tasks from queues and execute the actual business logic.
Hands-on Implementation
Let us implement an e-commerce checkout flow where we must:
-
Reserve inventory.
-
Charge the customer's credit card.
-
If the payment fails, release the reserved inventory.
1. Defining the Workflow and Saga Logic
Using the Temporal Go SDK, we can programmatically manage compensations using the built-in workflow.Saga helper.
package saga
import (
time "time"
"go.temporal.io/sdk/workflow"
)
type OrderRequest struct { OrderID string UserID string Amount float64 ItemID string Quantity int }
func CheckoutWorkflow(ctx workflow.Context, req OrderRequest) (err error) { // Set activity options with retry policies ao := workflow.ActivityOptions{ StartToCloseTimeout: 10 * time.Second, RetryPolicy: &temporal.RetryPolicy{ InitialInterval: time.Second, BackoffCoefficient: 2.0, MaximumAttempts: 5, }, } ctx = workflow.WithActivityOptions(ctx, ao)
// Initialize the Saga orchestrator
sagaOptions := &workflow.SagaOptions{
ParallelCompensation: false, // Run compensations sequentially
}
saga := &workflow.Saga{Options: sagaOptions}
// Defer execution of compensations in case of failure
defer func() {
if err != nil {
// Execute registered compensations in reverse order
sagaCtx, _ := workflow.NewDisconnectedContext(ctx)
_ = saga.Compensate(sagaCtx)
}
}()
// Step 1: Reserve Inventory
var inventoryResult string
err = workflow.ExecuteActivity(ctx, ReserveInventoryActivity, req.ItemID, req.Quantity).Get(ctx, &inventoryResult)
if err != nil {
return err
}
// Register compensation for Step 1
saga.AddCompensation(ReleaseInventoryActivity, req.ItemID, req.Quantity)
// Step 2: Authorize Payment
var paymentResult string
err = workflow.ExecuteActivity(ctx, AuthorizePaymentActivity, req.UserID, req.Amount).Get(ctx, &paymentResult)
if err != nil {
// If payment fails, deferred function kicks off ReleaseInventoryActivity
return err
}
return nil
}
2. Implementing Activities
Activities are standard Go functions. Because they run on independent workers, they should be designed as stateless wrappers around downstream RPC or DB calls.
package saga
import (
"context"
"fmt"
)
func ReserveInventoryActivity(ctx context.Context, itemID string, qty int) (string, error) { // Production code: Call external Inventory microservice fmt.Printf("Reserving %d of item %s\n", qty, itemID) return "Reserved Successfully", nil }
func ReleaseInventoryActivity(ctx context.Context, itemID string, qty int) (string, error) { // Compensation logic to revert reservation fmt.Printf("Releasing %d of item %s\n", qty, itemID) return "Released Successfully", nil }
func AuthorizePaymentActivity(ctx context.Context, userID string, amount float64) (string, error) { // Production code: Charge credit card via Stripe/Adyen if amount > 10000.0 { return "", fmt.Errorf("limit exceeded for user %s", userID) } fmt.Printf("Charged %.2f to user %s\n", amount, userID) return "Charged Successfully", nil }
Architectural Warning: Compensations must be designed to eventually succeed. Temporal will retry them indefinitely based on your configured retry policy, as leaving a compensation uncompleted results in inconsistent state across your system.
Enterprise Security Hardening & Best Practices
-
Idempotency Keys: Ensure all microservice API endpoints called by Activities are strictly idempotent. Since Temporal retries failed Activities, network blips can cause an activity to run multiple times. Pass a unique transaction key (e.g., the
WorkflowRunID) with every downstream call. -
Zero-Trust Payload Encryption: To protect sensitive customer information (such as credit card details and PII) passing through the Temporal Cluster, implement custom
DataConverterserialization. This encrypts payloads at the worker level, ensuring the Temporal Control Plane never sees plaintext sensitive data. -
Disconnected Context for Compensations: When executing compensations inside a deferred function, always wrap the execution context using
workflow.NewDisconnectedContext(ctx). This prevents compensations from being aborted if the parent workflow execution context is canceled or timed out. -
Split Task Queues: Separate high-throughput, low-latency activities (like internal calculations) from slow, external integrations (like payment gateway calls) into distinct task queues. This prevents external API latency from starving your worker thread pools.
Key Takeaways & Conclusion
Orchestrating distributed transactions no longer requires complex manual state tracking or brittle event choreography. By leveraging Temporal's stateful workflows with Go's performance capabilities, you can build self-healing, highly observable Sagas. If a microservice fails midway, your system reliably rolls back to a consistent state while preserving strict execution logs for auditing and debugging.
