Running Singleton Jobs Across Replicas¶
Every replica of a service runs the same code, so an in-process scheduler fires the same job on every replica. With five replicas, the nightly report is generated five times, the reminder emails go out five times, the cleanup runs five times in parallel and contends with itself. Leader election is one fix; for scheduled jobs a lighter one usually suffices: each replica claims the schedule slot — "rollup for 2026-10-02T03:00" — with an atomic set-if-absent, and only the claimant runs. Tested with 5 replicas firing a job for 10 consecutive slots against Redis 7.4, the job ran exactly 10 times — once per slot — and every other attempt saw the claim and skipped. The design question that remains is what happens when the claimant crashes mid-run, and this guide answers it explicitly.
Prerequisites¶
- Python 3.11+,
redis.asyncioand Redis (or Postgres for the database variant). - In-process scheduling, from scheduling cron jobs inside an asyncio service.
- Coordination options, from Distributed Locks & Coordination.
1. Derive a stable slot key¶
The claim only works if every replica computes the same key for the same run. Derive it from the scheduled time, never from "now", because replicas' clocks and scheduling jitter differ by milliseconds to seconds:
from datetime import datetime, timedelta, timezone
def slot_key(job: str, scheduled: datetime, period: timedelta) -> str:
epoch = datetime(1970, 1, 1, tzinfo=timezone.utc)
slot_start = epoch + ((scheduled - epoch) // period) * period
return f"ran:{job}:{slot_start.strftime('%Y-%m-%dT%H:%M')}"
slot_key("rollup", datetime(2026, 10, 2, 3, 0, 2, tzinfo=timezone.utc), timedelta(hours=1))
# 'ran:rollup:2026-10-02T03:00'
Flooring to the period means a replica that fires two seconds late still claims the same slot as one that fired on time. Use UTC throughout; a local-time key changes meaning twice a year at daylight-saving transitions, which is exactly when duplicate or skipped runs are hardest to notice.
Verify: compute the key on several replicas with skewed clocks; all produce the same string for the same scheduled run.
2. Claim the slot atomically, then run¶
The claim is a single SET … NX EX — the first replica to set it runs the job, everyone else skips:
async def run_once(r, job: str, scheduled, period, fn) -> bool:
key = slot_key(job, scheduled, period)
claimed = await r.set(key, replica_id(), nx=True, ex=int(period.total_seconds() * 3))
if not claimed:
return False # another replica has this slot
log.info("running %s for slot %s", job, key)
await fn()
return True
Measured across 5 replicas and 10 slots: 10 runs, 40 skips. The expiry only needs to outlast the window in which a late replica could still attempt the slot — a few periods is plenty — after which the key cleans itself up. Storing the replica id as the value makes "who ran it" answerable from Redis.
In that test the same replica won every slot, because all five attempted in the same order on one event loop; in production, scheduling jitter spreads claims naturally, and adding a small random delay before claiming spreads them deliberately if you want load to rotate.
Verify: with all replicas running, each slot's job runs once and the claiming replica is recorded.
3. Decide what a crash mid-run should mean¶
Claiming before running gives at-most-once: if the claimant crashes halfway, the slot stays claimed and the job does not run again for that slot. Claiming after would give at-least-once with duplicates. Neither is universally right:
| Job | Better on crash | Approach |
|---|---|---|
| send daily digest emails | skip rather than double-send | claim first (at-most-once) |
| nightly billing rollup | must complete, duplicates harmful | lease + completion marker + idempotency |
| cache refresh | either is fine | claim first |
| data export for a partner | must complete | lease + completion marker |
For jobs that must complete, combine a short lease (who is running it now) with a completion marker (it finished):
async def run_to_completion(r, job, scheduled, period, fn, lease_s: int = 60) -> bool:
base = slot_key(job, scheduled, period)
if await r.exists(base + ":done"):
return False # already completed by someone
if not await r.set(base + ":lease", replica_id(), nx=True, ex=lease_s):
return False # someone is running it now
await fn() # must be idempotent: a retry may repeat it
await r.set(base + ":done", replica_id(), ex=int(period.total_seconds() * 3))
return True
If the runner crashes, the lease expires and the next replica to fire — or a retry loop — runs it again; the done marker stops repeats after success. Because a re-run can repeat partial work, the job must be idempotent, using the techniques in making background jobs idempotent.
Verify: kill the runner mid-job; for the at-most-once variant the slot is skipped and logged, for the completion variant another replica completes it after the lease.
4. Retry missed slots without a scheduler of schedulers¶
With the completion variant, a crashed run is only retried if something fires the job again. Rather than waiting for the next scheduled tick, have each replica sweep recent slots on its own tick:
async def sweep(r, job, period, fn, lookback: int = 3) -> None:
now = datetime.now(timezone.utc)
for k in range(lookback, -1, -1): # oldest first
scheduled = now - k * period
await run_to_completion(r, job, scheduled, period, fn)
Each tick checks the current slot and the previous few; completed slots are skipped by their done marker, slots whose lease is held are skipped, and slots left unfinished by a crash are picked up. lookback limits how far back a job is worth catching up — a daily report from three days ago may not be. This turns the replicas themselves into the retry mechanism, with no extra scheduler process to keep alive.
Verify: crash a run, then let the next tick fire; the missed slot completes, and the current slot runs too.
5. Know when to use a job system instead¶
The claim pattern is right for a handful of periodic jobs inside a service. When jobs need retries with backoff, visibility into history, manual re-runs and alerting on failure, a job system with built-in uniqueness is less code to own — arq's cron jobs are unique across workers by default, as covered in running arq workers with Redis. The claim pattern can also enqueue rather than run: the winning replica enqueues one job for the slot, and the job system's workers handle execution and retries.
async def enqueue_once(r, arq, job, scheduled, period, *args) -> None:
if await r.set(slot_key(job, scheduled, period), replica_id(), nx=True,
ex=int(period.total_seconds() * 3)):
await arq.enqueue_job(job, *args, _job_id=slot_key(job, scheduled, period))
Passing the slot key as the job id makes the enqueue itself idempotent as well.
Verify: for each scheduled job, it is clear whether it runs in-process or through the job system, and duplicates are prevented at one layer.
Verification¶
Singleton scheduling works when:
- Every replica computes the same slot key, from the scheduled time in UTC.
- Each slot runs once, verified with all replicas active.
- Crash behaviour is chosen per job — skip or complete — and documented.
- Must-complete jobs are idempotent and retried by the sweep.
Diagnostic Hook: record, for each job and slot, which replica claimed it, when it started and when it finished, and export "slots without a done marker older than one period" as a gauge. A non-zero value is a run that crashed or hung; a slot claimed by no replica at all means the scheduler itself stopped firing, which is the failure nobody notices until the report is missing.
Pitfalls & edge cases¶
- Slot keys from wall-clock "now". Replicas a second apart can compute different slots and both run.
- Local time in slot keys. Daylight-saving changes duplicate or skip slots.
- Claim without expiry. Keys accumulate forever; expire after a few periods.
- Non-idempotent jobs with retries. A re-run after a crash repeats partial effects.
Frequently Asked Questions¶
How do I stop a cron job running on every replica?
Have each replica claim the schedule slot with an atomic set-if-absent, such as SET ran:job:slot NX EX in Redis, and run the job only if the claim succeeds. In testing, 5 replicas firing 10 slots produced exactly 10 runs.
What happens if the replica running a singleton job crashes?
With a claim taken before running, the slot is skipped: at-most-once. For jobs that must finish, use a short lease plus a completion marker so another replica retries after the lease expires, and make the job idempotent.
How should I build the key for a schedule slot?
From the scheduled time floored to the job's period, in UTC, plus the job name. Never from the current time, which differs slightly between replicas.
Should I use leader election or per-slot claims for scheduled jobs?
Per-slot claims for a few periodic jobs: they need no long-lived leader and no renewal. Leader election when one replica should own a continuous duty, such as a scheduler that enqueues many jobs.
Related¶
- Distributed Locks & Coordination — up to the topic overview.
- Leader election for asyncio workers — the alternative for continuous duties.
- Concurrent Execution & Worker Patterns — the section overview.