Hiring BarSupport

Managing Long-Running Processes — Cross-Cutting Pattern

Pattern35 min read5 diagrams

Technologies that implement this pattern: Kafka · PostgreSQL · DynamoDB · Redis · ZooKeeper & etcd · Flink & Stream Processing

Why This Matters#

Long-running processes are where the request/response model quietly stops being true. Every layer of a web stack assumes work finishes in seconds: API gateways cut requests at ~29–60s, load balancers drop idle connections at 60s, serverless functions die at 15 minutes, Kubernetes gives a terminating pod 30 seconds. Then the product asks for a 40-minute video export, a 3-day onboarding workflow with a human approval step, or a payout that waits for a bank settlement file. The work outlives every process, connection, and deploy that touches it.

Most candidates treat this as a queue question: "put it on a queue and have a worker process it." Staff engineers treat it as a durable state question: where does the progress of this work live, so that any worker can crash, any deploy can roll, and the work resumes from the last meaningful step instead of starting over or — worse — running twice? The queue moves work; it does not remember how far the work got. That memory — checkpoints, heartbeats, leases, compensation records — is the design.

In interviews, the pattern hides inside job schedulers, payment flows, report generation, video processing, ML training, data backfills, and order fulfillment. The L5 answer is "async job + queue + worker." The L6 answer adds the job state machine, heartbeats, checkpoint/resume, idempotent steps, cancellation, and a saga for multi-service work. The L7 answer asks whether the org should run one workflow platform instead of fifteen ad-hoc job tables, and what happens to 200K in-flight workflows when someone deploys incompatible code.

The 60-Second Version#

  • Anything over ~10 seconds leaves the request path. Return 202 Accepted with a job ID in < 200 ms; let the client poll (every 2–5 s with backoff) or receive a webhook/push when state changes.
  • Progress must live outside the worker. Checkpoint every 30–60 s or every N items to a durable store. A 2-hour job without checkpoints loses an average of 1 hour of work per crash.
  • Heartbeats detect death; leases prevent duplicates. Heartbeat every ~10–30 s; declare dead after 3 missed. SQS visibility timeout defaults to 30 s (max 12 h) — a job longer than its lease will run twice.
  • Every step must be idempotent. At-least-once delivery plus retries means each step runs ≥ 1 time. Key side effects on (job_id, step, attempt-independent key).
  • Multi-step, multi-service, or human-in-the-loop? Use a workflow engine. Temporal, AWS Step Functions (Standard: up to 1 year per execution), or similar give durable timers, retries, and history for free. Hand-rolling this is 12–18 engineer-months to reach parity.
  • Cancellation and compensation are features, not afterthoughts. A saga across 4 services needs a compensating action for each completed step, and cancellation must be cooperative — checked at every checkpoint boundary.

The Problem#

A user clicks "Export all my data," a merchant triggers a payout, an admin starts a 50M-row backfill, or a customer's order kicks off payment → inventory → shipping → notification across four services, one of which waits two days for a carrier pickup. Each takes minutes to days. During that time, workers are redeployed several times, a node dies, a downstream service has a 20-minute outage, and the user clicks "cancel" halfway. The system must report accurate progress, never lose work, never double-charge or double-ship, resume after crashes without restarting from zero, and cleanly undo partial work when asked to stop.


Case Studies That Use This Pattern#


The Core Tradeoff#

StrategyWhat WorksWhat BreaksWho Pays
Synchronous with long timeoutSimplest code; caller gets the answer directlyGateway/LB timeouts (29–60 s), pinned threads, retries duplicate workUsers see timeouts; API on-call sees thread-pool exhaustion
Async job: queue + worker + job tableDecouples caller, scales workers independently, simple to reason aboutNo resume without checkpoints; lease/visibility mistakes cause duplicates; multi-step flows become ad-hoc state machinesThe team that owns the job table and its DLQ — forever
Checkpointed batch jobResumes after crash, bounded reworkCheckpoint design is per-job; consistency between checkpoint and side effectsJob authors must reason about idempotency at every boundary
Durable workflow engine (Temporal, Step Functions, Cadence)Durable timers, retries, history, visibility, human-in-loop waits of daysDeterminism constraints, versioning in-flight workflows, event-history limits, a new platform to run or buyPlatform team (or vendor bill); developers learn a new programming model
Choreographed saga (events between services)No central coordinator, services stay autonomousNobody can answer "where is order 123?"; compensation logic scatteredOn-call debugging a stuck flow across 5 teams' logs
Orchestrated saga (central workflow)One place for state, retries, compensation, and visibilityOrchestrator is a critical dependency; coupling to its schemaOrchestrator owner is paged for every downstream failure

The L5 → L6 → L7 Contrast#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First move"Return a job ID and process it on a queue""Where does progress live? I'll define the job state machine and checkpoint boundaries first""How many teams run ad-hoc job tables today? This should be one workflow platform with a paved-road SDK"
FailureRetries the job on failureHeartbeats, leases shorter than job duration are extended, checkpoint/resume, DLQ with an ownerDesigns the org's posture for platform-wide failure: workflow engine outage, bad deploy hitting 200K in-flight executions
CorrectnessAssumes the job runs onceIdempotent steps keyed on stable IDs; exactly-once effect via dedupStandardizes idempotency and compensation contracts across services so sagas compose
Multi-serviceChains calls in one workerOrchestrated saga with compensation per step and cancellationDecides orchestration vs choreography as org policy, including who owns cross-team workflows
Change over timeDeploys new worker codeVersions workflow code; in-flight executions finish on old logicSets version-support windows and drain policy; long-lived workflows are a migration liability tracked like schema
CostNot mentionedSizes worker pools from arrival rate × durationPrices platform vs ad-hoc: engine cost, headcount, and the incidents ad-hoc systems cause
Why "First move" separates levels

"Job ID and a queue" is the correct skeleton, and an L5 who says it is not wrong. But a queue only answers who picks up the work, not how far it got. A worker that dies 90 minutes into a 2-hour export either restarts at zero (wasted work, doubled cost) or, worse, left half its side effects committed. Staff engineers begin with the state machine and checkpoint boundaries because those determine correctness; the queue is a transport detail. Principal engineers notice that every team solving this independently produces a slightly different, slightly broken state machine — and that the fix is a platform, not a better design review.

Why "Change over time" separates levels

A request lives for 200 ms, so deploying new code is invisible to it. A workflow lives for days or months — it will span dozens of deploys. If workflow code changes the order of steps while executions are mid-flight, replay-based engines like Temporal detect nondeterminism and the execution gets stuck; hand-rolled systems silently skip or repeat steps. Staff engineers version workflow logic explicitly. Principal engineers treat long-lived executions like data with a schema: they set a support window ("workflow code versions are supported for 30 days; longer workflows must use continue-as-new") and track in-flight executions per version as a migration metric.


Staff Default Position#

Get the work out of the request path, keep its progress in durable state, and make every step idempotent.

Return 202 with a job ID; persist a job record with an explicit state machine; process with workers that hold a lease, heartbeat, and checkpoint at meaningful boundaries; key every side effect on a stable idempotency key so retries are harmless. Use a plain queue + job table for single-step jobs under ~15 minutes. Use a durable workflow engine as soon as the work has multiple steps across services, waits longer than a worker can live (timers, human approvals, external callbacks), or needs compensation. Cancellation is cooperative and checked at every checkpoint.


When to Deviate#

  • Work reliably under ~5–10 seconds. Keep it synchronous with a timeout and idempotency key; async machinery adds polling latency and a state machine for no benefit.
  • Pure data processing at high volume. A 10 TB transformation is a batch/stream job (Spark, Flink), not 10 million workflow executions — engines like Temporal charge per action and aren't built for per-record orchestration.
  • One team, one simple job type, low volume. A Postgres job table with SELECT ... FOR UPDATE SKIP LOCKED is fine for thousands of jobs/day; adopting a workflow platform for one cron-like job is over-engineering.
  • Latency-critical multi-step flows. A workflow engine adds ~10–50 ms per step of persistence overhead. For a sub-100 ms checkout path, do the critical steps synchronously and hand only the tail (fulfillment, notifications) to the workflow.

Common Interview Mistakes#

What Candidates SayWhat Interviewers HearWhat Staff Engineers Say
"We'll increase the API timeout to 10 minutes""Every proxy and client in the path will cut this off""Anything over ~10 s returns 202 with a job ID; the client polls or gets a webhook."
"The worker processes the job from the queue""What happens when it dies halfway?""The worker holds a lease, heartbeats every 15 s, and checkpoints every 60 s. A new worker resumes from the last checkpoint."
"The job runs exactly once""I haven't thought about redelivery""Delivery is at-least-once. Every side effect uses an idempotency key so a second run is a no-op."
"If a step fails, we roll back""There's no distributed transaction to roll back""Completed steps are committed. I'll define a compensating action for each — refund for charge, release for reserve."
"We'll add a cancel flag""Who checks it, and what about the steps already done?""Cancellation is cooperative at checkpoint boundaries, and cancel triggers compensation for completed steps."
"We'll build a workflow system with a state table""You're about to rebuild Temporal badly""For multi-step flows with timers and compensation, I'd use a workflow engine; a job table is fine for single-step work."

Quick Reference#

Diagram: Quick Reference

Staff Sentence Templates#

"This takes [N minutes], which is longer than our gateway timeout and longer than a deploy cycle, so I'll return 202 with a job ID and make the job's progress durable — the worker is disposable, the checkpoint is not."

"Each worker holds a lease of [N] seconds and heartbeats every [N/3]. If it misses three heartbeats, another worker takes over from the last checkpoint. Because delivery is at-least-once, every side effect is keyed on [job_id + step]."

"This flow crosses [N] services and waits on [a human / an external callback] for up to [N days], so I'd model it as a workflow in [Temporal / Step Functions] rather than a chain of queue consumers — we get durable timers, retries, and a single place to answer 'where is order X?'"

"If the user cancels after [step K], steps 1 through K are already committed. Cancellation runs the compensations in reverse — [refund, release inventory, void label] — and each compensation is itself idempotent and retried until it succeeds."


Implementation Deep Dive#

1. The Async Job Pattern — 202, Job Resource, and Status Delivery#

The API contract is the first design decision, because clients will build around it for years.

Diagram: 1. The Async Job Pattern — 202, Job Resource, and Status Delivery
# Job resource — the contract clients poll
GET /jobs/j_81
{
  "id": "j_81",
  "state": "RUNNING",              # QUEUED | RUNNING | SUCCEEDED | FAILED | CANCELLING | CANCELLED
  "progress": { "done": 4_200_000, "total": 10_000_000 },
  "created_at": "...", "updated_at": "...",
  "result_url": null,              # set on SUCCEEDED; short-lived signed URL
  "error": null,                   # { code, message, retryable } on FAILED
  "links": { "cancel": "/jobs/j_81/cancel" }
}
# Response header: Retry-After: 5   -> server controls polling rate
Delivery MethodLatencyCostUse When
Polling with Retry-After2–10 s1 cheap read per pollDefault; works everywhere, server throttles clients
Webhook (signed, retried with backoff for ~24–72 h)< 1 sDelivery infra + retry queueB2B integrations; client has a server
Push (WebSocket/SSE/mobile push)< 1 sConnection infraInteractive UIs already holding a connection
Long-poll GET /jobs/:id?wait=30s< 1 sHeld connection per waiterSimple clients wanting low latency without webhooks

🎯 Staff Insight: The Idempotency-Key on job creation is not optional. A user double-clicking "Export" or a mobile client retrying on a flaky network will otherwise create two 40-minute jobs. A unique constraint on (tenant_id, idempotency_key) turns the retry into "here's the job you already started."

2. Leases and Heartbeats — Knowing When a Worker Is Dead#

A worker that stops making progress looks identical to a worker that's slow — until you add heartbeats. The lease is the mechanism that lets someone else take over, and it is also the mechanism that causes duplicate execution if you get it wrong.

LEASE_TTL = 60s
HEARTBEAT_EVERY = 15s          # ~1/4 of lease: tolerates 3 missed beats

function run(job):
    lease = acquire(job.id, owner=worker_id, ttl=LEASE_TTL)   # CAS on job row
    stop = start_heartbeat(job.id, lease, every=HEARTBEAT_EVERY)
    try:
        for chunk in job.remaining_chunks(from=job.checkpoint):
            if stop.lease_lost:                  # heartbeat CAS failed: someone took over
                return abandon()                 # stop immediately, do NOT commit
            process(chunk)                       # idempotent
            checkpoint(job.id, chunk.end, fencing=lease.token)
    finally:
        stop()

function checkpoint(job_id, cursor, fencing):
    # Fencing token rejects writes from a worker whose lease already expired
    ok = db.exec("UPDATE jobs SET cursor=$1, updated_at=now()
                  WHERE id=$2 AND lease_token=$3", cursor, job_id, fencing)
    if not ok: raise LeaseLost

Queue-native equivalents: SQS visibility timeout (default 30 s, max 12 h) extended with ChangeMessageVisibility as the heartbeat; Temporal activity heartbeat() with a HeartbeatTimeout; Kubernetes Jobs with activeDeadlineSeconds as a hard cap.

ParameterDefault I'd PickToo ShortToo Long
Lease TTL60 sGC pause or slow I/O triggers takeover → duplicate workDead worker's job stalls for the full TTL
Heartbeat intervalTTL / 4Heartbeat load on the job store (10K workers × 1/s = 10K writes/s)Fewer chances before expiry
Hard deadline per job3× p99 durationKills legitimate slow jobsZombie jobs hold capacity for hours

🎯 Staff Insight: A lease without a fencing token is only a hint. If a paused worker wakes up after its lease expired, it will happily write its stale checkpoint over the new owner's progress. The fencing token (a monotonically increasing lease version checked on every write) is what makes takeover safe.

3. Checkpointing — Bounding the Cost of a Crash#

Checkpoint frequency is a direct tradeoff between overhead and rework. The expected rework per crash is roughly half the checkpoint interval.

# Cursor-based checkpoint for a 50M-row backfill
function backfill(job):
    cursor = job.checkpoint or MIN_ID
    while cursor < job.max_id:
        batch = db.query("SELECT ... WHERE id > $1 ORDER BY id LIMIT 5000", cursor)
        write_idempotently(batch)                 # upsert keyed on id: replay-safe
        cursor = batch.last_id
        if time_since_last_checkpoint() > 30s or rows_since_checkpoint() > 50_000:
            checkpoint(job.id, cursor)
        throttle_to(db.replica_lag < 2s)          # don't melt the primary
Checkpoint IntervalOverhead (write ~5 ms)Expected Rework per CrashUse For
Every itemHigh — dominates for small items~0Expensive, non-idempotent items (payments)
Every 30–60 s< 0.1%15–30 sDefault for most batch jobs
Every 10 minNegligible~5 minCheap-to-redo compute
Never0Half the jobOnly jobs under ~1–2 min

The consistency trap: the checkpoint and the side effects are two writes. If you write side effects and crash before checkpointing, the replay repeats those side effects — which is fine only if they're idempotent. If you checkpoint first and crash before side effects, you skip work. Rule: side effects first, idempotently; checkpoint after. Or commit both in the same transaction when they share a database.

🎯 Staff Insight: Checkpoint meaning, not just position. A cursor says "I've processed up to row 4.2M." It doesn't say whether the aggregate you're building includes rows 4.1–4.2M. For stateful jobs, the checkpoint must include the partial state (or its location) — this is exactly what Flink's checkpoint barriers do for stream jobs.

4. Durable Workflow Engines — Temporal and Step Functions#

When work spans services, timers, and humans, a workflow engine persists every step's result in an event history. Workflow code is replayed from that history after a crash, so the code looks sequential but survives process death.

# Temporal-style workflow (pseudocode): order fulfillment
@workflow
def fulfill_order(order):
    opts = ActivityOptions(start_to_close=30s, retry=Retry(max_attempts=10,
                           backoff=2.0, max_interval=5min))
    charge = activity(charge_card, order, opts)           # idempotency key = order.id
    try:
        activity(reserve_inventory, order, opts)
        label = activity(create_shipping_label, order, opts)
        # Durable timer: survives worker restarts and deploys, costs nothing while waiting
        picked_up = wait_for_signal("carrier_pickup", timeout=72h)
        if not picked_up:
            activity(escalate_to_ops, order, opts)
            wait_for_signal("carrier_pickup", timeout=7d)
        activity(send_shipped_email, order, opts)
    except (CancelledError, ActivityFailure) as e:
        # Compensation in reverse order; each is idempotent and retried
        activity(release_inventory, order, opts)
        activity(refund_charge, charge, opts)
        raise

# Rules the engine imposes: workflow code must be deterministic —
# no wall-clock reads, random(), direct I/O, or unversioned logic changes.
EngineMax ExecutionHistory LimitPricing ShapeBest For
Temporal (self-hosted or Cloud)Unbounded (use continue-as-new)~50K events / 50 MB per execution (warns ~10K)Self-hosted: ~1–2 FTE ops; Cloud: per actionCode-first workflows, complex branching, polyglot teams
AWS Step Functions Standard1 year25,000 history events~$0.025 per 1,000 state transitionsAWS-native orchestration of Lambda/ECS/services
AWS Step Functions Express5 minutesN/A (at-least-once)Per request + durationHigh-volume, short event processing
Hand-rolled job table + queueWhatever you buildWhatever you build12–18 engineer-months + an owner foreverSingle-step jobs only
Diagram: 4. Durable Workflow Engines — Temporal and Step Functions

🎯 Staff Insight: The event history limit is the one production teams learn the hard way. A workflow that loops over 100K items as activities will hit ~50K events and be terminated. Long-lived or high-iteration workflows must use continue-as-new (restart with carried-over state) or fan out to child workflows — design for it on day one.

5. Sagas, Cancellation, and Compensation#

A multi-service long-running process has no global transaction. Each step commits locally; failure after step K means undoing steps 1..K with semantic compensations — which are business operations, not database rollbacks.

Forward StepCompensationCompensation Can Fail BecauseHandling
Charge cardRefund (not void after capture)Processor outageRetry with backoff up to 72 h; alert finance at 24 h
Reserve inventoryRelease reservationService downReservation TTL (e.g., 30 min) is the safety net
Create shipping labelVoid labelCarrier already scanned itNot compensable → escalate to human ops queue
Send confirmation emailSend correction email—Some actions can only be countered, never undone
function cancel(job_id, reason):
    db.cas(job_id, from=["QUEUED","RUNNING","WAITING"], to="CANCELLING", reason)
    # Cooperative: worker/workflow observes CANCELLING at its next checkpoint or await
    # QUEUED jobs are cancelled immediately — no work to undo

# Worker side, at every checkpoint boundary:
if job.state == "CANCELLING":
    run_compensations(completed_steps.reversed())
    db.cas(job.id, from="CANCELLING", to="CANCELLED")
    return

Pivot step: order the saga so that irreversible steps (shipping, external transfers) come last, after all compensable steps succeed. Once past the pivot, the saga can only move forward — retries, not compensation.

🎯 Staff Insight: Cancellation latency is a product decision. "Cancel stops within one checkpoint interval (≤ 60 s)" is an SLA you can state and meet. "Cancel stops immediately" is a promise you can't keep for a worker mid-way through a 5-minute external call. Say the bound out loud.


Architecture Diagram#

Diagram: Architecture Diagram

Reading the diagram: two execution paths share one control plane. Single-step jobs go queue → worker with leases and checkpoints in the job store. Multi-service flows go to the workflow engine, which owns timers and compensation. The reaper and DLQ are the correctness backstops — every stuck or poisoned job ends up somewhere a human owns.


Failure Scenarios#

1. Visibility Timeout Shorter Than the Job — Eight Copies of Every Report#

Context: A reporting service uses SQS with the default 30-second visibility timeout. Reports used to take ~10 s. A new enterprise customer's reports take 4–6 minutes. Workers don't extend visibility.

t=0        Enterprise customer schedules 2,000 monthly reports at 09:00
t=+30s     Visibility expires on in-progress messages; SQS redelivers to other workers
t=+2min    Each report now running on 3-5 workers simultaneously
t=+5min    Worker pool (200) fully consumed by duplicates; other tenants' jobs queue
t=+8min    Reporting DB replica CPU 98%; reports emailed 3-8 times each
t=+25min   On-call raises visibility to 15 min by config; duplicates stop after drain

Detection: jobs.duplicate_start (same job_id started > 1 time), sqs.ApproximateReceiveCount p99 > 2, worker.utilization pinned at 100% while jobs.completed_rate falls.

Blast radius: All tenants on the shared worker pool; customer-visible duplicate emails; reporting replica saturation.

Mitigation: Raise visibility timeout; purge duplicate in-flight work by checking a job_id lease before starting.

Prevention: (1) Heartbeat with ChangeMessageVisibility every 1/4 of the timeout; (2) lease + fencing on the job row so a second receive can't start the job; (3) dedup email sends on (job_id, recipient); (4) per-tenant concurrency caps so one tenant's large jobs can't consume the pool.

Owner: Reporting team owns the worker; the jobs platform team owns the SDK default that should have heartbeated automatically.

🎯 Staff Insight: Duration is a property of the workload, not the code — it changes when customers change. Any lease or timeout set as a constant will eventually be shorter than some job. Heartbeat-extended leases make duration irrelevant; fixed timeouts make it a ticking bomb.

2. The Nondeterministic Deploy — 40K Workflows Stuck Mid-Flight#

Context: An order workflow in a replay-based engine is modified to add a fraud check between charge and inventory reservation. The change is deployed without versioning. ~40K orders are in-flight, many waiting on 72-hour carrier-pickup timers.

t=0        Deploy rolls out new workflow code to all workers
t=+1min    Workers replay in-flight histories; new code schedules fraud_check where
           history recorded reserve_inventory -> nondeterminism error
t=+5min    Every in-flight workflow that wakes (signal, timer) gets stuck; new orders fine
t=+40min   Ops notices stuck-workflow count climbing; shipping confirmations stop
t=+1h      Rollback to old code; stuck workflows resume on replay
t=+1h30    New code re-deployed behind a version gate: old executions keep old path

Detection: workflow.task_failures{type="nondeterminism"} > 0 (should always be 0), workflow.stuck_count by workflow type, and business signal orders.shipped_rate dropping while orders.created_rate is flat.

Blast radius: Every in-flight execution of that workflow type — orders placed over the previous 3 days — while new orders looked healthy, masking the incident.

Mitigation: Roll back worker code; replay recovers without data loss because history is intact.

Prevention: (1) Version gates (if version >= 2: fraud_check()) for any change to step order; (2) CI replay tests that run new code against a sample of recorded production histories; (3) a deploy check that blocks if in-flight executions exist on an unsupported version.

Owner: Order team for the workflow code; workflow platform team for the replay-test tooling and deploy gate.

🎯 Staff Insight: Long-lived workflows turn every code change into a schema migration against in-flight data. The question at review isn't "does the new code work?" — it's "does the new code correctly continue an execution that started under the old code three days ago?"

3. The Compensation That Never Ran — Charged but Not Shipped#

Context: A choreographed saga: OrderPlaced → PaymentCaptured → InventoryReserved → ShipmentCreated. The inventory service fails permanently for a discontinued SKU and publishes InventoryFailed. The payment service's consumer for that event was removed during a refactor six weeks earlier.

Week 0       Refactor removes payment-service handler for InventoryFailed
Week 0-6     ~1,900 orders: charged, never reserved, never shipped, never refunded
Week 6       Chargeback rate for the merchant category doubles; finance escalates
Week 6+2d    Reconciliation query finds charged-but-unshipped orders older than 7 days
Week 6+5d    Manual refunds issued; saga moved to an orchestrator with explicit compensation

Detection (should have been): A reconciliation invariant — count(orders WHERE captured AND NOT shipped AND age > 72h AND NOT refunded) — alerting at > 0; saga.incomplete_age_p99 per flow.

Blast radius: ~1,900 customers charged for goods never shipped; chargeback fees and a payments-risk review.

Mitigation: Batch refunds with idempotency keys; customer outreach.

Prevention: Orchestrated saga where compensation is code in one place, not an event subscription someone can delete; end-to-end invariant checks run hourly; contract tests that every failure event has a registered handler.

Owner: Checkout team owns the orchestrated saga; finance ops owns the reconciliation invariant.

🧭 Principal Insight: Choreography distributes the happy path nicely and distributes the failure path dangerously — no one team owns "what happens when step 3 fails." For flows that move money or goods, the org should mandate orchestration with a named owner of the end-to-end state.

Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Lease shorter than jobjobs.duplicate_start, receive count p99Shared worker pool, duplicate side effectsHeartbeat-extended leases, fencingJob owner + jobs platform (SDK default)
Worker crash without checkpointjobs.rework_seconds, restart countThat job's duration × crashesCheckpoint every 30–60 sJob owner
Nondeterministic workflow deployworkflow.task_failures{nondeterminism}All in-flight executions of that typeRoll back; version gates; replay testsWorkflow owner + platform
History limit hitworkflow.history_events near 50KExecutions terminatedContinue-as-new, child workflowsWorkflow owner
Missing / failed compensationSaga invariant queries, saga.incomplete_age_p99Money or goods in inconsistent stateOrchestrated compensation with retries; human queueFlow owner + finance ops
Workflow engine outageengine.persistence_latency_p99, task backlogEvery workflow in the org pausesWorkers retry; engine HA across AZs; runbookWorkflow platform team
Zombie jobs holding capacityjobs.running_age > hard deadlineWorker pool capacityReaper enforces deadlines, cancelsJobs platform
DLQ accumulationdlq.depth > 0 for 10 minSilent failed workAlert, triage, redrive toolingOwning team of each queue

The Principal Lens#

Why L7 Sees This Problem Differently#

A Staff engineer designs one durable, idempotent, cancellable job system. A Principal engineer walks the org and finds fifteen: a jobs table in billing, a Redis-backed queue in growth, a cron-plus-status-column in data, Celery in ML, Step Functions in one AWS-heavy team, and an in-house "workflow framework" whose author left. Each has its own lease bugs, its own DLQ nobody watches, and its own answer to "where is this customer's request?" The L7 insight is that durable execution is infrastructure, like a database — it should be provided once, operated by a team that specializes in it, and consumed through an SDK that makes heartbeats, idempotency, and versioning the default rather than the exception.

The Org-Level Fault Line#

One durable-execution platform vs fit-for-purpose tools per team. A single platform (e.g., Temporal) gives consistent semantics, visibility across flows, and one on-call for the hard parts. But it concentrates risk — an engine outage pauses every workflow in the company — and it imposes a programming model (determinism, versioning) that some workloads don't need. The Principal position: one paved-road workflow engine for multi-step, cross-service, or long-waiting flows; one paved-road job queue SDK for single-step async work; batch/stream engines for bulk data. Three tools with clear boundaries beats one tool stretched over everything, and beats fifteen tools chosen by accident.

Cost Model#

ScaleWorkloadAd-Hoc per TeamManaged EngineSelf-Hosted EngineAssumptions
Small (5 teams)~1M workflow actions/mo~0.5 FTE per team maintaining job tables (~2.5 FTE ≈ $750K/yr)~$1–3K/mo engine + 0.25 FTENot worth it: ~1.5 FTE ≈ $450K/yrManaged wins by a wide margin
Mid (40 teams)~500M actions/mo~10–15 FTE scattered + incidents~$30–100K/mo + 2 FTE platform~$15–30K/mo infra + 4–5 FTEManaged vs self-host is close; decide on data residency and team capacity
Large (300 teams)~20B actions/moUnmanageable; multiple Sev1s/yr from bespoke systems~$1M+/mo at list~$150–300K/mo infra + 10–15 FTESelf-hosted or negotiated enterprise contract; engine becomes tier-0

Assumptions: $300K fully loaded FTE; managed pricing at list order of magnitude per action (Step Functions Standard ~$25 per million transitions; hosted Temporal pricing is per action in a similar order); persistence on Cassandra/Postgres for self-hosted. The biggest saving is not infra — it's retiring 10+ bespoke job systems and their incidents.

🧭 Principal Move: "I don't want a workflow engine because it's elegant. I want it because we currently have fifteen homemade ones, three of which caused Sev1s last year. The business case is headcount returned to product teams and incidents that stop happening."

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionReversibilityCost to ReverseWhy
Polling vs webhook for job statusTwo-wayAdd the other in ~2–4 engineer-weeksAdditive API change
Queue technology for single-step jobsTwo-way~1 engineer-month behind an SDKJobs are short-lived; drain and switch
Checkpoint formatTwo-way-ishMust support old format until in-flight jobs finishBounded by max job duration
Workflow engine choiceOne-wayQuarters: every workflow rewritten to a new model; in-flight executions must drain (days to months)Programming model is embedded in code, not config
Orchestration vs choreography for money flowsOne-way-ishRe-architecting across 4–6 teamsOwnership and event contracts spread across services
Public job API shape (states, IDs)One-wayClients integrate once and never updateExternal contracts live for years

The Standard I'd Write#

RFC: Long-Running and Asynchronous Work

Scope: Any operation that can exceed 10 seconds, spans more than one service, or waits on timers, humans, or external callbacks.

Mandatory requirements:

  • Operations that can exceed 10 s MUST return 202 with a job resource; synchronous endpoints MUST NOT exceed the gateway timeout.
  • Job creation MUST accept an idempotency key; every side-effecting step MUST be idempotent on a stable key.
  • Workers MUST use heartbeat-extended leases with fencing tokens via the paved-road job SDK; fixed visibility timeouts are prohibited for jobs over 30 s.
  • Jobs over 2 minutes MUST checkpoint at least every 60 s.
  • Flows that move money or goods across services MUST use the orchestrated workflow engine with an explicit compensation for every compensable step and a documented pivot step.
  • Workflow code changes MUST pass replay tests against recorded production histories; step-order changes MUST be version-gated.
  • Every queue and DLQ MUST have an owning team and a depth alert.

Exceptions: Bulk data processing uses batch/stream platforms instead; exceptions approved by the workflow platform lead, reviewed every 12 months.

Success metrics: Bespoke job systems retired (target: 15 → 2 in 24 months); jobs.duplicate_start rate < 0.01%; zero nondeterminism incidents reaching production; saga.incomplete_age_p99 under SLA for every money flow.

What I'd Tell the VP#

We have about fifteen different homemade systems for running background work, and three of them caused major incidents last year — duplicate charges, stuck orders, and reports sent eight times. I want to replace them with one supported system for multi-step business processes and one standard library for simple background jobs, run by a small platform team. Product teams stop maintaining plumbing, and we get a single place to answer "where is this customer's order?" The cost is a platform team of four to five engineers and a vendor or infrastructure bill that is smaller than the engineering time we already spend on these systems.

Principal Interview Signals#

SignalWhat It Sounds Like
Durable execution as infrastructure"Workflows are infrastructure like a database. Teams shouldn't each build leases and retries any more than they'd each build a storage engine."
Clear tool boundaries"Engine for multi-step flows, job SDK for single steps, Flink or Spark for bulk data — three tools with clear lines."
Change management for in-flight state"Long-lived workflows make every deploy a migration. I'd set a version-support window and gate deploys on replay tests."
Concentration risk"Once 60% of tier-1 flows run on the engine, it's tier-0. It needs multi-AZ persistence, a game day, and a runbook for 'engine is down.'"
Ownership of failure paths"For money flows, someone must own the end-to-end failure path. That's why I'd mandate orchestration over choreography there."

Staff answers that L7 interviewers find insufficient:

  • "I'd use Temporal for this flow" — right for one flow, silent on the other fourteen job systems and on who runs the engine.
  • "We'll add heartbeats and checkpoints" — correct mechanics, but no mechanism to make every team do it by default (an SDK, a standard, a lint).
  • "Compensation handles failures" — without asking who owns compensations that span teams, or how the org detects the one that silently never ran.

In the Wild#

Uber: Cadence, and the Birth of Temporal#

Uber built Cadence, an open-source durable workflow engine, to run long-lived business processes — the kind that previously lived in ad-hoc combinations of queues, cron jobs, and state tables across many teams. Workflow code is written as ordinary sequential code; the engine persists an event history and replays it to rebuild state after failures. Cadence's creators later founded Temporal, which forked and evolved the project into what is now one of the most widely adopted durable-execution engines.

Staff insight: The origin story is the argument. A large engineering org found that every team was independently building fragile versions of the same thing — retries, timers, state persistence — and the fix was a platform with a programming model, not better code reviews. In an interview, citing Cadence/Temporal signals you understand why workflow engines exist, not just that they do.

Netflix: Conductor and the Orchestration Choice#

Netflix built and open-sourced Conductor (2016), a JSON-defined orchestration engine, to coordinate microservice workflows such as content ingestion and encoding pipelines. Netflix's public writing on it explicitly favored central orchestration over pure event choreography, citing visibility and control over complex flows. In late 2023 Netflix announced it would stop maintaining the open-source project; community and commercial forks continue it.

Staff insight: Two lessons in one. First, orchestration won inside a company famous for microservice autonomy because someone must be able to see the whole flow. Second, even a strong in-house engine became a maintenance burden its creator chose to hand off — the build-vs-buy question applies to workflow engines too.

AWS Step Functions: Workflows as a Managed Service#

AWS Step Functions (launched 2016) offers state-machine workflows defined in the Amazon States Language. Standard workflows can run for up to one year with exactly-once step semantics and full execution history; Express workflows run up to five minutes at much higher throughput with at-least-once semantics. Pricing for Standard is per state transition.

Staff insight: The Standard/Express split is the whole pattern in one product decision: long-lived, auditable, per-step-billed durability vs short, high-volume, cheap execution. Mentioning which one you'd pick — and that 10M items as individual Standard transitions is an expensive mistake — shows you think about cost shape, not just features.


Practice Drill#

Prompt: "Design 'Export my account data' for a SaaS product. Exports can be 50 GB and take up to 2 hours. Users should see progress, be able to cancel, and get a download link when done."

Staff Answer

POST /exports with an Idempotency-Key returns 202 and a job resource in < 200 ms; a unique constraint on (account_id, idempotency_key) makes double-clicks safe, and a per-account limit of one active export prevents abuse. The job is a single logical process with many chunks, so a job SDK is enough — no workflow engine. A worker acquires a 60 s lease with a fencing token and heartbeats every 15 s; it walks each data category by cursor, streams rows into multipart-upload parts in object storage, and checkpoints {category, cursor, upload_id, parts[]} every 60 s after the part is durably uploaded, so a crash loses at most a minute and replays are idempotent (part numbers are deterministic from the cursor). Progress is rows done / estimated total, surfaced via GET /jobs/:id with Retry-After: 5, plus an email when done. Cancel sets CANCELLING; the worker sees it at the next checkpoint (≤ 60 s), aborts the multipart upload, and marks CANCELLED. On success, the result is a signed download URL valid for 24 h; the object expires after 7 days via lifecycle rule. A reaper fails jobs with no heartbeat for 3 minutes and re-queues them from checkpoint; after 3 failures, the job goes to a DLQ the owning team is alerted on. Reads come from a replica, throttled on replica lag < 2 s.

Why this is L6:

  • Picks the lighter tool (job SDK) because it's one logical step, and says why
  • Defines checkpoint content and ordering (side effect first, checkpoint after)
  • States a cancellation latency bound and cleans up partial artifacts
  • Protects the primary database and names the DLQ owner

What L7 adds:

  • Makes data export a platform capability (GDPR/CCPA access requests use the same pipeline) with a compliance SLA, e.g., 30 days, and audit logs
  • Standardizes the job resource shape across all async APIs so clients integrate once
  • Prices it: exports at 1% of accounts/month × 50 GB drives replica and egress capacity planning

Prompt: "An order workflow has been 'processing' for 9 days for 3,000 customers. You're called in. What do you do?"

Staff Answer

First, classify: are they stuck (engine errors, nondeterminism, exhausted retries) or legitimately waiting (a timer or signal that never came)? Pull workflow.task_failures by type and a sample of histories. If a recent deploy correlates, roll back and version-gate. If they're waiting on a signal (e.g., carrier pickup) that was never sent, find the broken producer, replay the signals idempotently, and add a timeout branch so no wait is unbounded. For customers, send status updates and decide with the business whether to compensate (refund) or push forward. Then add the missing guardrails: every wait has a timeout with an escalation path, workflow.oldest_open_age per type alerts at 2× expected duration, and replay tests gate deploys.

Why this is L6:

  • Separates "stuck" from "waiting forever" before acting
  • Fixes the individual executions and the class of failure (unbounded waits)

What L7 adds:

  • Makes "every wait has a timeout and an owner" a platform lint, not a code-review hope
  • Reviews which other flows share the same broken signal producer — correlated risk across teams

Staff Interview Application#

How to Introduce This Pattern#

"This operation takes longer than any request timeout and longer than a deploy cycle, so I'll take it out of the request path and focus on where its progress lives. The API returns a job ID; workers hold heartbeat-extended leases and checkpoint durably so any crash resumes, not restarts; every side effect is idempotent because delivery is at-least-once. If the flow spans services or waits on humans or timers, I'd put it in a workflow engine and define compensations for each step."

Lead with durable progress, not with the queue — that's the sentence that signals Staff.

When NOT to Use This Pattern#

  • Work that reliably finishes in a few seconds. Keep it synchronous; a job resource and polling add latency and complexity for nothing.
  • High-volume, per-record data transformation. Use Spark/Flink/batch jobs with their own checkpointing, not one workflow per record.
  • Fire-and-forget best-effort work. Analytics pings or cache warmups can go on a queue without leases, checkpoints, or sagas — losing some is acceptable.
  • A single team's low-volume job. A Postgres table with SKIP LOCKED beats adopting and operating a workflow platform.

Follow-Up Questions to Anticipate#

Interviewer AsksWhat They Are TestingHow to Respond
"What if the worker dies halfway?"Durability of progress"Heartbeats stop, the lease expires in ~60 s, another worker resumes from the last checkpoint — at most ~60 s of rework."
"How do you avoid doing the work twice?"Idempotency and fencing"Delivery is at-least-once, so effects are keyed on job and step. Fencing tokens reject writes from a worker whose lease expired."
"How does cancel work?"Cooperative cancellation"Cancel sets CANCELLING; the worker checks at every checkpoint, runs compensations in reverse, then marks CANCELLED. Bound: one checkpoint interval."
"Why not just a queue?"Queue vs workflow judgment"A queue moves work; it doesn't remember progress, timers, or compensation. For single-step jobs a queue plus job table is enough; for multi-step sagas I'd use an engine."
"How do you deploy changes to a 3-day workflow?"Versioning in-flight state"Version-gate step-order changes so old executions continue on old logic, and run replay tests against recorded histories in CI."
"What if a compensation fails?"Saga failure depth"Compensations retry with backoff for hours; non-compensable steps go to a human queue. I'd order the saga so irreversible steps come after the pivot."

Evaluation Rubric#

DimensionSenior (L5)Staff (L6)Principal (L7)
API contractAsync with job ID202 + job resource, idempotency key, Retry-After, webhooksOne job-resource standard across all async APIs
DurabilityRetry on failureLeases, heartbeats, fencing, checkpoints, reaperPaved-road SDK makes these defaults; bespoke systems retired
Multi-step correctnessSequential calls in a workerOrchestrated saga, compensation, pivot step, cancellation boundOrg policy: orchestration for money/goods flows, owners for failure paths
EvolutionRedeploy workersVersioned workflows, replay testsVersion-support windows, engine as tier-0 with game days

Strong hire signals: asks where progress lives before drawing a queue; sets lease/heartbeat/checkpoint numbers; designs cancellation and compensation unprompted; knows the engine's constraints (determinism, history limits).

Lean no-hire signals: raises the HTTP timeout; assumes exactly-once execution; "rollback" across services; no DLQ owner.

Common false positive: Fluency in a specific engine's API ≠ long-running-process judgment. A candidate can write perfect Temporal code and still put a 100K-iteration loop in one workflow, skip versioning, and choreograph a money flow nobody owns.


Capacity Planning Quick Reference#

NumberValueContext
Move off the request path> ~10 sGateways cut at ~29–60 s; mobile clients give up sooner
Lease TTL / heartbeat60 s / 15 s3 missed heartbeats before takeover
Checkpoint interval30–60 sExpected rework ≈ half the interval
SQS visibility timeout30 s default, 12 h maxExtend via heartbeat; never rely on the default
Serverless function max (AWS Lambda)15 minAnything longer needs a job or workflow
Pod termination grace (Kubernetes)30 s defaultCheckpoint on SIGTERM; don't start new chunks
Step Functions Standard≤ 1 year, 25K history events~$0.025 per 1,000 transitions
Step Functions Express≤ 5 minHigh-volume, at-least-once
Temporal history limit~50K events / 50 MBContinue-as-new well before; warns ~10K
Workflow engine step overhead~10–50 ms per persisted stepKeep latency-critical steps synchronous
Worker pool sizingarrival rate × avg duration × 1.310 jobs/s × 120 s ≈ 1,200 slots + 30% headroom
Webhook retry window24–72 h with backoffSign payloads; include event ID for dedup

Pitfalls checklist:

  • Is anything over ~10 s still synchronous?
  • Where does progress live if the worker dies right now?
  • Is the lease heartbeat-extended and fenced, or a fixed timeout?
  • Is every side effect idempotent on a stable key?
  • Does every wait have a timeout and an escalation path?
  • Does every compensable step have a compensation, and are irreversible steps after the pivot?
  • How do in-flight executions survive a code change?
  • Who owns the DLQ, and what alerts on it?
  1. Loading the index…