Software - General
1867048 Members
1269 Online
110506 Solutions
New Discussion

Temporal: Building Long-Running Systems

 
Prithviraj-
HPE Pro

Temporal: Building Long-Running Systems

Most backend systems start simple. A request comes in, you do some work, you return a response. Then reality shows up:

  • A payment must be charged, then a receipt emailed, then inventory reserved — and each step can fail independently.
  • An order must wait 3 days for the customer to confirm, and auto-cancel otherwise.
  • A data pipeline must process 10,000 records, resume from record 6,832 after a crash, and never double-process a record.

The usual answer is a pile of cron jobs, message queues, a status column in a database, and retry loops. It works until a process gets OOM-killed halfway through step 3, and you spend the next two hours reconstructing what happened from logs.

Temporal exists to remove that entire category of work.


1. What Temporal actually is

Temporal is a durable execution platform. You write your business logic as ordinary code — loops, conditionals, function calls, sleep — and Temporal guarantees that this code runs to completion exactly once, even if the process running it crashes, the machine dies, or the datacenter goes away for an hour.

The core trick: Temporal does not persist your program's memory. It persists the history of everything that happened. When your process dies and a new one picks up the work, Temporal replays that history through your function to rebuild its exact state, then continues from where it left off.

Two things follow from this, and almost everything else in Temporal is a consequence of them:

  1. Your workflow code must be deterministic, because it gets replayed.
  2. Anything non-deterministic (network calls, DB reads, random numbers, clock reads) must be moved out into activities, whose results are recorded in history and returned verbatim on replay.

2. The building blocks Workflow

The orchestrator. It decides what happens and in what order. It is deterministic and must not perform I/O.

package app

import (
	"time"

	"go.temporal.io/sdk/temporal"
	"go.temporal.io/sdk/workflow"
)

type OrderInput struct {
	OrderID string
	UserID  string
	Amount  int64
}

func OrderWorkflow(ctx workflow.Context, in OrderInput) (string, error) {
	ao := workflow.ActivityOptions{
		StartToCloseTimeout: 30 * time.Second,
		RetryPolicy: &temporal.RetryPolicy{
			InitialInterval:    time.Second,
			BackoffCoefficient: 2.0,
			MaximumInterval:    time.Minute,
			MaximumAttempts:    5,
		},
	}
	ctx = workflow.WithActivityOptions(ctx, ao)

	var chargeID string
	if err := workflow.ExecuteActivity(ctx, ChargeCard, in.UserID, in.Amount).Get(ctx, &chargeID); err != nil {
		return "", err
	}

	if err := workflow.ExecuteActivity(ctx, ReserveInventory, in.OrderID).Get(ctx, nil); err != nil {
		// compensate: the money is already taken
		_ = workflow.ExecuteActivity(ctx, RefundCharge, chargeID).Get(ctx, nil)
		return "", err
	}

	if err := workflow.ExecuteActivity(ctx, SendReceipt, in.UserID, chargeID).Get(ctx, nil); err != nil {
		return "", err
	}

	return chargeID, nil
}

Read that again: it looks like a plain Go function. There is no state machine, no status column, no retry loop, no resume logic. If the worker running this dies between ChargeCard and ReserveInventory, another worker picks it up, replays history, sees ChargeCard already returned chargeID, and proceeds straight to ReserveInventory. The card is not charged twice.

Activity

The part that touches the outside world. Activities may be non-deterministic, may fail, and are automatically retried according to their retry policy.

package app

import (
	"context"
	"fmt"
)

func ChargeCard(ctx context.Context, userID string, amount int64) (string, error) {
	resp, err := paymentGateway.Charge(ctx, userID, amount)
	if err != nil {
		return "", fmt.Errorf("charge failed: %w", err)
	}
	return resp.ChargeID, nil
}

func ReserveInventory(ctx context.Context, orderID string) error {
	return inventory.Reserve(ctx, orderID)
}

Activities receive a normal context.Context. It is cancelled when the activity is cancelled or its timeout fires, so pass it down to your HTTP and DB calls.

Worker

The process that hosts your code and polls Temporal for work. It is just a long-running Go binary.

package main

import (
	"log"

	"go.temporal.io/sdk/client"
	"go.temporal.io/sdk/worker"

	"example.com/app"
)

func main() {
	c, err := client.Dial(client.Options{
		HostPort:  client.DefaultHostPort, // localhost:7233
		Namespace: "default",
	})
	if err != nil {
		log.Fatalln("dial:", err)
	}
	defer c.Close()

	w := worker.New(c, "orders", worker.Options{})

	w.RegisterWorkflow(app.OrderWorkflow)
	w.RegisterActivity(app.ChargeCard)
	w.RegisterActivity(app.ReserveInventory)
	w.RegisterActivity(app.RefundCharge)
	w.RegisterActivity(app.SendReceipt)

	if err := w.Run(worker.InterruptCh()); err != nil {
		log.Fatalln("worker:", err)
	}
}

Starting a workflow

opts := client.StartWorkflowOptions{
	ID:        "order-" + in.OrderID, // your idempotency key
	TaskQueue: "orders",
}

run, err := c.ExecuteWorkflow(ctx, opts, app.OrderWorkflow, in)
if err != nil {
	return err
}

var chargeID string
err = run.Get(ctx, &chargeID) // blocks until the workflow completes

The workflow ID is your deduplication mechanism. Calling ExecuteWorkflow twice with the same ID while the first is still running returns a handle to the existing run rather than starting a second one. Derive it from a business key — order ID, user ID, invoice number — never from a UUID generated at call time.


3. Task queues and namespaces

Task queue is the routing mechanism. Workers poll a named queue; workflows and activities are dispatched to a queue. That's it — there is no broker configuration, no exchange, no topic.

This gives you useful control for free:

// Route heavy GPU/CPU work to a dedicated fleet
ctx = workflow.WithActivityOptions(ctx, workflow.ActivityOptions{
	TaskQueue:           "video-encoding",
	StartToCloseTimeout: 2 * time.Hour,
})

Namespace is the isolation boundary — think of it as a virtual cluster. Use separate namespaces for prod, staging, and per-team workloads. Retention, archival, and access controls are configured per namespace.


4. Timeouts: the four that matter

This is where most people get burned, so be explicit about all of them.

Timeout Meaning Guidance ScheduleToStartTimeout How long a task may sit in the queue before a worker picks it up Usually leave unset; a non-zero value here is a capacity alarm, not a correctness control StartToCloseTimeout Max duration of a single attempt Always set this. Make it slightly larger than your realistic p99 ScheduleToCloseTimeout Total wall-clock budget across all retries Set when there is a real business deadline HeartbeatTimeout Max gap between heartbeats for long activities Required for anything that runs more than ~a minute

And on the workflow itself:

opts := client.StartWorkflowOptions{
	ID:                       "order-123",
	TaskQueue:                "orders",
	WorkflowExecutionTimeout: 7 * 24 * time.Hour, // total lifetime, retries included
	WorkflowRunTimeout:       24 * time.Hour,     // single run
	WorkflowTaskTimeout:      10 * time.Second,   // one decision step
}

5. Retries and error semantics

Activities retry by default (infinitely, with exponential backoff, unless you say otherwise). Workflows do not retry by default — they are already durable, so a workflow "failing" means your business logic decided to fail.

Some errors should never be retried. A 400 Bad Request will still be a 400 on the fifth attempt.

// Mark an error as non-retryable from inside an activity
return temporal.NewNonRetryableApplicationError(
	"card declined", "CardDeclined", nil,
)
// Or exclude by type in the retry policy
RetryPolicy: &temporal.RetryPolicy{
	MaximumAttempts:        5,
	NonRetryableErrorTypes: []string{"CardDeclined", "InvalidInput"},
}

Handling failures in the workflow:

var appErr *temporal.ApplicationError
err := workflow.ExecuteActivity(ctx, ChargeCard, in.UserID, in.Amount).Get(ctx, &chargeID)
if errors.As(err, &appErr) && appErr.Type() == "CardDeclined" {
	return workflow.ExecuteActivity(ctx, NotifyDeclined, in.UserID).Get(ctx, nil)
}

Retries make at-least-once the default. Design activities to be idempotent. Pass an idempotency key (ActivityInfo.WorkflowExecution.ID + activity ID is a good one) to any external API that supports it.


6. Heartbeats and long-running activities

If an activity takes minutes or hours, heartbeat so Temporal can detect a dead worker quickly instead of waiting out the full StartToCloseTimeout. Heartbeats can also carry progress, which lets a retry resume instead of restart.

func ProcessRecords(ctx context.Context, batchID string) error {
	start := 0
	if activity.HasHeartbeatDetails(ctx) {
		_ = activity.GetHeartbeatDetails(ctx, &start) // resume where we died
	}

	records, err := loadRecords(ctx, batchID)
	if err != nil {
		return err
	}

	for i := start; i < len(records); i++ {
		if err := handle(ctx, records[i]); err != nil {
			return err
		}
		activity.RecordHeartbeat(ctx, i+1)

		select {
		case <-ctx.Done():
			return ctx.Err() // cancelled or timed out
		default:
		}
	}
	return nil
}

Set HeartbeatTimeout in the activity options or the heartbeats are ignored.


7. Timers: sleeping for days is normal

// Wait three days. Costs nothing while waiting — no thread, no memory, no cron.
if err := workflow.Sleep(ctx, 72*time.Hour); err != nil {
	return err
}

A sleeping workflow is not resident in any process. Temporal stores a timer and re-dispatches the workflow when it fires. Sleeping for a year is as cheap as sleeping for a second.

Never use time.Sleep, time.Now(), rand, os.Getenv, or map iteration order in workflow code. Use the deterministic equivalents:

now := workflow.Now(ctx)
id := workflow.GetInfo(ctx).WorkflowExecution.ID

var uuid string
_ = workflow.SideEffect(ctx, func(ctx workflow.Context) interface{} {
	return newUUID()
}).Get(&uuid)

8. Signals: sending data into a running workflow

Signals are how the outside world talks to a workflow in flight.

func SubscriptionWorkflow(ctx workflow.Context, userID string) error {
	cancelCh := workflow.GetSignalChannel(ctx, "cancel")
	cancelled := false

	for month := 0; month < 12; month++ {
		timerCtx, cancelTimer := workflow.WithCancel(ctx)
		timer := workflow.NewTimer(timerCtx, 30*24*time.Hour)

		sel := workflow.NewSelector(ctx)
		sel.AddFuture(timer, func(workflow.Future) {})
		sel.AddReceive(cancelCh, func(c workflow.ReceiveChannel, _ bool) {
			c.Receive(ctx, nil)
			cancelled = true
			cancelTimer()
		})
		sel.Select(ctx)

		if cancelled {
			return workflow.ExecuteActivity(ctx, SendCancellationEmail, userID).Get(ctx, nil)
		}
		if err := workflow.ExecuteActivity(ctx, ChargeMonthly, userID).Get(ctx, nil); err != nil {
			return err
		}
	}
	return nil
}

From a client:

err := c.SignalWorkflow(ctx, "subscription-user-42", "", "cancel", nil)

SignalWithStartWorkflow starts the workflow if it isn't running and signals it if it is — the standard pattern for "append to a session that may or may not exist yet".

Human-in-the-loop falls out of this naturally: workflow.Await on an approval signal with a timeout, and escalate if nobody responds.

approved := false
approvalCh := workflow.GetSignalChannel(ctx, "approval")
workflow.Go(ctx, func(ctx workflow.Context) {
	var decision bool
	approvalCh.Receive(ctx, &decision)
	approved = decision
})

ok, _ := workflow.AwaitWithTimeout(ctx, 48*time.Hour, func() bool { return approved })
if !ok {
	return workflow.ExecuteActivity(ctx, EscalateToManager, in).Get(ctx, nil)
}

Queries: reading state out

Queries are read-only, synchronous, and must not mutate state or schedule work.

func OrderWorkflow(ctx workflow.Context, in OrderInput) error {
	status := "pending"
	if err := workflow.SetQueryHandler(ctx, "status", func() (string, error) {
		return status, nil
	}); err != nil {
		return err
	}
	// ... status = "charged" ... status = "shipped" ...
	return nil
}
resp, _ := c.QueryWorkflow(ctx, "order-123", "", "status")
var status string
_ = resp.Get(&status)

Updates: signal + query in one round trip

An update sends data in and returns a result, with optional validation that rejects bad input without writing to history.

_ = workflow.SetUpdateHandlerWithOptions(ctx, "addItem",
	func(ctx workflow.Context, item Item) (int, error) {
		cart = append(cart, item)
		return len(cart), nil
	},
	workflow.UpdateHandlerOptions{
		Validator: func(ctx workflow.Context, item Item) error {
			if item.Qty <= 0 {
				return errors.New("quantity must be positive")
			}
			return nil
		},
	},
)

9. Child workflows and fan-out

Use a child workflow when a sub-process deserves its own ID, history, retry policy, and lifecycle. Use an activity when it's just a unit of work.

// Sequential child
var res ShipmentResult
err := workflow.ExecuteChildWorkflow(ctx, ShipmentWorkflow, orderID).Get(ctx, &res)
// Parallel fan-out over activities
futures := make([]workflow.Future, 0, len(itemIDs))
for _, id := range itemIDs {
	futures = append(futures, workflow.ExecuteActivity(ctx, ReserveItem, id))
}
for _, f := range futures {
	if err := f.Get(ctx, nil); err != nil {
		return err
	}
}

ParentClosePolicy controls what happens to children when the parent finishes — terminate them (default), abandon them, or request cancellation.


10. Continue-As-New: keeping history bounded

Event history is not infinite. Practically, keep it under a few thousand events. For workflows that loop forever (subscriptions, monitors, per-entity actors), atomically restart with fresh history:

func MonitorWorkflow(ctx workflow.Context, state State) error {
	for i := 0; i < 500; i++ {
		if err := workflow.Sleep(ctx, time.Hour); err != nil {
			return err
		}
		if err := workflow.ExecuteActivity(ctx, Poll, state.Target).Get(ctx, &state); err != nil {
			return err
		}
	}
	return workflow.NewContinueAsNewError(ctx, MonitorWorkflow, state)
}

The workflow ID stays the same; only the run ID changes. Signals received while continuing-as-new are carried over.


11. Determinism and versioning

Because history is replayed, you cannot freely change workflow code that has in-flight executions. Reordering activities, adding a step in the middle, or removing a Sleep will make replay diverge from history and produce a non-determinism error.

Two safe options:

Patching — for incremental changes:

if workflow.GetVersion(ctx, "add-fraud-check", workflow.DefaultVersion, 1) == workflow.DefaultVersion {
	// old path: existing runs continue as before
} else {
	// new path: only new runs take this
	if err := workflow.ExecuteActivity(ctx, FraudCheck, in).Get(ctx, nil); err != nil {
		return err
	}
}

Worker Versioning / new workflow type — for large rewrites: pin old runs to old workers and route new runs to the new build.

Safe changes that never need a patch: activity implementation bodies, timeout and retry policy values, and log statements.


12. Testing

Temporal's Go SDK ships a test framework that runs workflows in-process with a skipping clock — a workflow that sleeps 30 days finishes in milliseconds.

func TestOrderWorkflow(t *testing.T) {
	var s testsuite.WorkflowTestSuite
	env := s.NewTestWorkflowEnvironment()

	env.OnActivity(ChargeCard, mock.Anything, "u1", int64(500)).Return("ch_1", nil)
	env.OnActivity(ReserveInventory, mock.Anything, "o1").Return(nil)
	env.OnActivity(SendReceipt, mock.Anything, "u1", "ch_1").Return(nil)

	env.ExecuteWorkflow(OrderWorkflow, OrderInput{OrderID: "o1", UserID: "u1", Amount: 500})

	require.True(t, env.IsWorkflowCompleted())
	require.NoError(t, env.GetWorkflowError())

	var chargeID string
	require.NoError(t, env.GetWorkflowResult(&chargeID))
	require.Equal(t, "ch_1", chargeID)
	env.AssertExpectations(t)
}

Test compensation paths by returning errors from mocks, and test signals with env.RegisterDelayedCallback.

Replay tests are your safety net against accidental non-determinism: download the history of a production workflow and replay it against your new code in CI.

replayer := worker.NewWorkflowReplayer()
replayer.RegisterWorkflow(OrderWorkflow)
err := replayer.ReplayWorkflowHistoryFromJSONFile(nil, "testdata/order_history.json")
require.NoError(t, err)

13. Schedules and cron

Recurring work is first-class — no external cron, no missed-fire mystery.

_, err := c.ScheduleClient().Create(ctx, client.ScheduleOptions{
	ID: "nightly-reconciliation",
	Spec: client.ScheduleSpec{
		CronExpressions: []string{"0 2 * * *"},
		TimeZoneName:    "Asia/Kolkata",
	},
	Action: &client.ScheduleWorkflowAction{
		ID:        "reconcile",
		Workflow:  ReconcileWorkflow,
		TaskQueue: "batch",
	},
	Overlap: enumspb.SCHEDULE_OVERLAP_POLICY_SKIP,
})

Schedules can be paused, backfilled, and triggered manually. The overlap policy answers "what if the previous run is still going?" — a question plain cron never asks.


14. Search attributes and visibility

Beyond fetching a workflow by ID, you can index and query workflows by business fields.

_ = workflow.UpsertSearchAttributes(ctx, map[string]interface{}{
	"CustomerId":  in.UserID,
	"OrderStatus": "charged",
})
resp, _ := c.ListWorkflow(ctx, &workflowservice.ListWorkflowExecutionsRequest{
	Query: `WorkflowType = 'OrderWorkflow' AND OrderStatus = 'charged' AND ExecutionStatus = 'Running'`,
})

This is how you answer "which orders are stuck in payment right now?" without adding a reporting table.


15. Cancellation and cleanup

Cancellation is cooperative and propagates to running activities and child workflows. Cleanup must run in a disconnected context, since the normal one is already cancelled.

defer func() {
	if errors.Is(ctx.Err(), workflow.ErrCanceled) {
		newCtx, _ := workflow.NewDisconnectedContext(ctx)
		_ = workflow.ExecuteActivity(newCtx, ReleaseInventory, in.OrderID).Get(newCtx, nil)
	}
}()

Note the difference: cancel is graceful and lets the workflow run its cleanup; terminate is immediate and runs nothing. Prefer cancel.


16. Running it locally

brew install temporal          # or: https://temporal.io/setup/install-temporal-cli
temporal server start-dev      # server on :7233, Web UI on http://localhost:8233
temporal workflow start \
  --task-queue orders \
  --type OrderWorkflow \
  --workflow-id order-123 \
  --input '{"OrderID":"o1","UserID":"u1","Amount":500}'

temporal workflow show   --workflow-id order-123
temporal workflow signal --workflow-id order-123 --name cancel
temporal workflow query  --workflow-id order-123 --type status

The Web UI shows the full event history of every execution — every activity attempt, every failure, every retry, with inputs and outputs. This alone replaces a large amount of custom logging and admin tooling.


17. Production notes

  • Payload size. Inputs, outputs, and signals go into history. Keep them small (rule of thumb: under ~100 KB, hard limit 2 MB). Pass S3 keys or row IDs, not blobs. Use a Data Converter for encryption or compression.
  • Worker tuning. MaxConcurrentActivityExecutionSize and MaxConcurrentWorkflowTaskExecutionSize are your throughput knobs. Run workflow and activity workers separately when activity work is heavy, so a saturated activity pool never starves workflow progress.
  • Retention. Closed workflow histories are deleted after the namespace retention period (default 3 days on the dev server, commonly 30 days in production). Ship anything you need long-term to your own store.
  • Metrics. The SDK exports Prometheus metrics. Watch temporal_workflow_task_schedule_to_start_latency (workers under-provisioned) and temporal_activity_execution_failed.
  • Sticky execution. Workers cache workflow state so they don't replay full history on every task. A cache miss just means a replay — correct, but slower. Size WorkflowCacheSize accordingly.

18. When not to use Temporal

Temporal is not free. It adds a server (plus Cassandra/PostgreSQL/MySQL and Elasticsearch), a programming model with real constraints, and per-step latency measured in milliseconds rather than microseconds.

Skip it when:

  • The work is a single, fast, idempotent operation. Just do it inline.
  • You need sub-millisecond latency per step.
  • You need high-throughput stream processing — that's Kafka/Flink territory, not Temporal's.
  • The whole system is one service with one database and no multi-step failure modes worth modelling.

Reach for it when:

  • A process spans multiple services or multiple external APIs.
  • A process spans minutes, days, or months.
  • Partial failure leaves the system in a state someone has to fix by hand.
  • You need auditability of what happened, step by step.
  • You've already written a status column, a retry table, and a reconciliation cron for the same flow.

19. The mental model to keep

Temporal turns a distributed, failure-prone, multi-step process into a single function you can read top to bottom.

Everything else — the event history, replay, task queues, timers, heartbeats — is machinery in service of that one property. Once you internalise "workflows orchestrate and must be deterministic; activities do the work and may fail", the rest of the API stops feeling like magic and starts feeling obvious.

Start with one flow. Pick the one that currently has a status column, a retry cron, and a runbook. Port it. You'll delete more code than you write.



I work at HPE
HPE Support Center offers support for your HPE services and products when and how you need it. Get started with HPE Support Center today.
[Any personal opinions expressed are mine, and not official statements on behalf of Hewlett Packard Enterprise]
Accept or Kudo
1 REPLY 1
Thaufique_Mod
Community Manager

Re: Temporal: Building Long-Running Systems

Thanks for sharing this informative topic. It is very helpful.



I work at HPE
HPE Support Center offers support for your HPE services and products when and how you need it. Get started with HPE Support Center today.
[Any personal opinions expressed are mine, and not official statements on behalf of Hewlett Packard Enterprise]
Accept or Kudo