When to Move from Job Queues to Temporal

Durable workflow engines such as Temporal remove whole categories of hand-built machinery — saga logs, sweepers, timer queues — at the price of a new piece of infrastructure and a new programming model. This guide, part of Job Chaining & Workflow Orchestration in Queue Fundamentals & Architecture, lays out the concrete signals that a queue-based workflow has outgrown its queue, shows the same workflow in both styles, and describes how to migrate one workflow at a time without a flag day.

Problem Statement

A subscription onboarding flow runs on Sidekiq: create the account, provision resources, send a welcome email, wait three days, check whether the customer has activated, and either send a nudge or hand the account to sales. It is implemented as eight job classes, a onboarding_runs table, a delayed-job for the three-day wait, and a nightly sweeper that fixes runs whose delayed job went missing. Every change touches four files, a bug last quarter sent 1,200 duplicate nudges after a Redis failover, and nobody can say with confidence which runs are waiting on what. The team is asking whether a workflow engine would fix this or just move the problem.

Prerequisites

  • An inventory of the workflows you run on queues today: steps, waits, human inputs, compensation, typical duration.
  • Numbers for each: runs per day, p99 duration, and how many incidents or manual fixes it caused in the last quarter.
  • A realistic view of who would operate a new stateful cluster, or budget for a managed offering (Temporal Cloud, AWS Step Functions, Azure Durable Functions).

Step 1 — Score Your Workflows Against the Signals

A queue is excellent at "do this unit of work reliably, soon". It is a poor substrate for "remember where this long-running process is and continue it later". The signals below each indicate that you are building workflow-engine features by hand. Score each workflow — two or more strong signals usually justify the move for that workflow.

Signal What you built by hand Engine feature that replaces it
Waits longer than minutes (days, weeks) Delayed jobs + a sweeper for lost timers Durable timers (workflow.sleep)
Waits on an external event or human Polling jobs, status columns, callbacks Signals and updates
Compensation across several steps Saga log + compensation jobs Try/except with compensation in code
Per-step timeouts and retries differ Per-job retry config spread across classes Activity options per call
"Where is run X?" is hard to answer Workflow tables, admin pages Built-in history and UI per run
Versioning pain with in-flight runs Keeping old job classes around Workflow versioning / patching APIs

A workflow that is three steps, completes in seconds, and never waits scores zero. Keep it on the queue; an engine would add latency and moving parts for no gain. The chain-and-chord patterns in Celery chains, groups, and chords are the right tool there.

Where each tool fits A chart with workflow duration on the horizontal axis and number of waits or human steps on the vertical. Short pipelines with no waits, such as thumbnail generation, sit in the queue region at the lower left. Onboarding flows with multi-day waits and approvals sit in the workflow-engine region at the upper right. Duration and waits decide the tool run duration: seconds to weeks waits job queue thumbnails, emails, ETL chunks workflow engine onboarding, approvals, sagas grey zone: score it

Step 2 — See the Same Workflow in Both Styles

The queue version of the onboarding flow spreads one process across jobs, a table, and a scheduler. Here is its core, condensed:

# Sidekiq version (condensed): state in a table, waits as delayed jobs
class ProvisionJob
  include Sidekiq::Job
  def perform(run_id)
    run = OnboardingRun.find(run_id)
    return if run.step_done?(:provision)
    Provisioner.call(run.account_id, idempotency_key: "#{run_id}:provision")
    run.mark!(:provision)
    WelcomeEmailJob.perform_async(run_id)
  end
end

class WelcomeEmailJob
  include Sidekiq::Job
  def perform(run_id)
    run = OnboardingRun.find(run_id)
    Mailer.welcome(run.account_id).deliver_now unless run.step_done?(:welcome)
    run.mark!(:welcome)
    ActivationCheckJob.perform_in(3.days, run_id)   # the timer lives in Redis
  end
end
# ...plus ActivationCheckJob, NudgeJob, HandoffJob, and a nightly sweeper

The Temporal version expresses the same process as ordinary code. Each execute_activity call is a durable step; the three-day sleep is a durable timer stored by the Temporal server, not in Redis.

# Temporal (Python SDK) version: the process is the code
from datetime import timedelta
from temporalio import workflow
from temporalio.common import RetryPolicy

with workflow.unsafe.imports_passed_through():
    from activities import provision, send_welcome, is_activated, send_nudge, hand_to_sales

@workflow.defn
class Onboarding:
    @workflow.run
    async def run(self, account_id: str) -> str:
        opts = dict(start_to_close_timeout=timedelta(minutes=5),
                    retry_policy=RetryPolicy(maximum_attempts=10, backoff_coefficient=2.0))
        await workflow.execute_activity(provision, account_id, **opts)
        await workflow.execute_activity(send_welcome, account_id, **opts)
        await workflow.sleep(timedelta(days=3))                  # survives restarts and deploys
        if await workflow.execute_activity(is_activated, account_id, **opts):
            return "activated"
        await workflow.execute_activity(send_nudge, account_id, **opts)
        await workflow.sleep(timedelta(days=4))
        if not await workflow.execute_activity(is_activated, account_id, **opts):
            await workflow.execute_activity(hand_to_sales, account_id, **opts)
        return "handed_off"

The engine persists an event history for each run. If a worker dies mid-run, another worker replays the history to rebuild the workflow's state and continues from the last completed step. That replay is also the constraint: workflow code must be deterministic (no wall-clock reads, random numbers, or direct I/O outside activities).

Where the workflow state lives In the queue version, the current step is in a database table, pending timers are delayed jobs in Redis, and a nightly sweeper repairs runs when the two disagree. In the engine version, one event history per run records every completed step and pending timer, and workers replay it to resume. Three places vs one queue-based runs table: step Redis: delayed jobs sweeper reconciles when they disagree workflow engine event history per run completed activities, timers, signals replayed to resume after a crash

Step 3 — Count the Real Costs

The engine removes code, but it adds a system. Be honest about both sides before committing.

# cost-comparison.yaml — fill in for your workflow
queue_based:
  code: 8 job classes, 1 table, 1 sweeper        # ~900 lines
  incidents_last_quarter: 3                       # duplicate nudges, lost timers
  infra: existing Redis + Sidekiq
temporal:
  code: 1 workflow, 5 activities                  # ~250 lines
  new_infra: Temporal server + Cassandra/Postgres + Elasticsearch (visibility)
             or Temporal Cloud (priced per action)
  new_skills: determinism rules, versioning with workflow.patched()
  latency: +10-50 ms per activity dispatch vs a direct enqueue
  migration: dual-run for one full workflow duration (7 days here)

Self-hosting Temporal means operating a persistence store, the frontend/history/matching services, and upgrades — a platform team's job. If there is no such team, the managed service or a cloud-native alternative (Step Functions on AWS) is usually the realistic choice. For workflows measured in seconds with no waits, the per-activity overhead is pure cost.

Step 4 — Migrate One Workflow, New Runs Only

Never move in-flight runs. Start new runs on the engine and let old runs finish on the queue; after one maximum workflow duration, delete the old code.

# Router: new onboarding runs go to Temporal behind a flag; old runs are untouched
class StartOnboarding
  def self.call(account_id)
    if Flags.enabled?(:onboarding_on_temporal, account_id)
      TemporalClient.start_workflow("Onboarding", account_id,
        id: "onboarding-#{account_id}",          # workflow id = natural key: no duplicates
        task_queue: "onboarding")
    else
      run = OnboardingRun.create!(account_id: account_id)
      ProvisionJob.perform_async(run.id)
    end
  end
end

Using a natural key as the workflow id gives you deduplication for free: starting a second workflow with the same id while the first is running fails, which is exactly the protection the Sidekiq version was missing when it sent duplicate nudges. Activities can reuse the existing service objects (Provisioner, Mailer) unchanged, which keeps the port small.

Step 5 — Map Queue Concepts onto Engine Concepts

Porting goes faster when the team shares a vocabulary for what each existing piece becomes. Most of the translation is mechanical:

On the job queue In the workflow engine What changes
A job class that calls an API An activity Same code; retries configured per call site
Enqueueing the next job await execute_activity(...) Sequencing becomes ordinary control flow
perform_in(3.days) / countdown workflow.sleep(timedelta(days=3)) Timer stored durably by the server
Status column + polling job Signal or update handler External events delivered to the running workflow
Saga log + compensation jobs try/except with compensations Compensation order lives in one function
Sweeper for stuck runs Workflow and activity timeouts Enforced by the engine, surfaced in its UI
Queue name per worker pool Task queue Same routing idea, same worker-pool separation

Two rows need more thought than the rest. Signals replace status polling: instead of a job that checks every hour whether a customer has activated, the application sends an activated signal to the workflow when it happens, and the workflow waits for either the signal or a timer, whichever comes first. Task queues in the engine play the same role as named queues in a job framework — route CPU-heavy activities to a separate worker pool exactly as you would with Celery task routing.

@workflow.defn
class Onboarding:
    def __init__(self) -> None:
        self.activated = False

    @workflow.signal
    def mark_activated(self) -> None:
        self.activated = True

    @workflow.run
    async def run(self, account_id: str) -> str:
        # ...provision and welcome as before...
        try:
            await workflow.wait_condition(lambda: self.activated, timeout=timedelta(days=3))
            return "activated"
        except asyncio.TimeoutError:
            await workflow.execute_activity(send_nudge, account_id,
                                            start_to_close_timeout=timedelta(minutes=5))
            return "nudged"
Polling becomes waiting On the queue, a job wakes every hour for three days to check a status column, running 72 times. In the engine, the workflow waits on a condition with a three-day timeout; the application sends a signal when the customer activates, and the workflow wakes once, either on the signal or on the timer. 72 polls vs one wake-up queue hourly check job, 72 runs, each a DB query engine wait_condition(activated, timeout = 3 days) signal arrives: wake once The waiting workflow consumes no worker time; the server holds the timer.

Verification

Before ramping the flag past a small percentage, confirm three things:

# 1. Runs complete and match the queue version's outcomes
temporal workflow list --query 'WorkflowType="Onboarding" AND ExecutionStatus="Completed"' | wc -l

# 2. Timers survive a worker restart: start a run with a short test sleep, kill workers, restart
kubectl rollout restart deployment/onboarding-worker
temporal workflow describe --workflow-id onboarding-acct-test-1   # still Running, timer pending

# 3. Duplicate starts are rejected
temporal workflow start --type Onboarding --task-queue onboarding \
  --workflow-id onboarding-acct-test-1 --input '"acct-test-1"'   # expect: already started

Compare business outcomes — activation rate, nudge counts — between the flagged and unflagged cohorts for a full cycle. A difference means the port changed behaviour, not just plumbing.

Gotchas & Edge Cases

Non-deterministic workflow code. Calling datetime.now() or reading a feature flag inside workflow code breaks replay. Use workflow.now() and pass flags in as inputs or read them in activities.

Changing a running workflow's code. Editing the sequence of activities breaks replay of in-flight runs. Use the SDK's patching/versioning API for every change to workflow logic, and keep old branches until old runs finish.

Large payloads in history. Activity inputs and outputs are stored in the history. Pass ids and storage keys, as you would through a queue — the claim-check pattern applies unchanged.

Using the engine as a queue. A workflow per event for millions of tiny, independent jobs is expensive and slow compared to a queue. Keep high-volume, fire-and-forget work on the queue and let workflows enqueue into it where needed.

FAQ

Is Temporal a replacement for Celery or Sidekiq? For long-running, stateful processes, largely yes. For high-volume, short, independent jobs, no — a queue is simpler, cheaper, and faster. Most organisations that adopt an engine keep their job queue.

What about AWS Step Functions? Step Functions is a managed workflow engine defined in JSON/YAML state machines rather than code. It has no servers to run and integrates tightly with AWS services, at the cost of a less expressive programming model and per-transition pricing. The decision signals in Step 1 apply equally.

Can I get most of the benefit without an engine? Often. A workflow table, idempotent steps, and a deadline sweeper — the approach in implementing sagas with compensating jobs — handle many workflows well. The engine wins when waits are long and workflows are numerous or change often.

Related