Distributed Cron Scheduling — Complete Guide
In this tutorial, you will learn about Distributed Cron Scheduling. We cover key concepts, practical examples, and best practices to help you master this topic.
Implement distributed cron across multiple servers using leader election, Redis distributed locks, and scheduler coordination to prevent duplicate execution in multi-node environments.
What You Learn
You will learn how to run Cron Jobs across multiple servers without duplicate execution using leader election, distributed locks with Redis, consensus-based scheduling, and handling failover when a scheduler node goes down.
Why It Matters
In a multi-server environment, every server runs the same crontab. Without coordination, backup runs on all servers simultaneously, thrashing the database. Distributed cron ensures each job runs exactly once across the entire cluster.
Real-World Use
DodaTech runs 10 application servers, each with identical crontabs. A Redis-based leader election ensures only one server runs each cron job. If the leader fails, another server takes over within 30 seconds. No duplicates, no missed schedules.
Simple Leader Election
#!/usr/bin/env python3
"""Leader election for distributed cron using Redis."""
import redis
import time
import uuid
import sys
import os
r = redis.Redis()
class LeaderElection:
def __init__(self, leader_key='cron:leader', ttl=30):
self.leader_key = leader_key
self.ttl = ttl
self.node_id = f"{os.uname().nodename}-{uuid.uuid4().hex[:6]}"
self.is_leader = False
def try_become_leader(self):
acquired = r.setnx(self.leader_key, self.node_id)
if acquired:
r.expire(self.leader_key, self.ttl)
self.is_leader = True
print(f"{self.node_id} became leader")
return True
current_leader = r.get(self.leader_key)
if current_leader:
current_leader = current_leader.decode()
print(f"Leader is {current_leader} (we are {self.node_id})")
self.is_leader = (current_leader == self.node_id)
return self.is_leader
def renew_leadership(self):
if self.is_leader:
current = r.get(self.leader_key)
if current and current.decode() == self.node_id:
r.expire(self.leader_key, self.ttl)
return True
self.is_leader = False
return False
def step_down(self):
if self.is_leader:
current = r.get(self.leader_key)
if current and current.decode() == self.node_id:
r.delete(self.leader_key)
self.is_leader = False
print(f"{self.node_id} stepped down")
leader = LeaderElection()
if leader.try_become_leader():
print("Running leader-only cron job")
time.sleep(5)
leader.step_down()
else:
print("Not the leader, skipping")
Expected output (first run):
server-01-a1b2c3 became leader
Running leader-only cron job
server-01-a1b2c3 stepped down
Distributed Cron Scheduler
#!/usr/bin/env python3
"""Full distributed cron scheduler with leader election."""
import redis
import time
import json
import hashlib
from datetime import datetime
r = redis.Redis()
class DistributedCronScheduler:
def __init__(self, node_id, ttl=30):
self.node_id = node_id
self.ttl = ttl
self.leader_key = 'distributed_cron:leader'
self.schedule_key = 'distributed_cron:schedules'
self.execution_key = 'distributed_cron:executed'
def register_schedule(self, name, cron_expr, task):
schedule = {
'name': name,
'cron': cron_expr,
'task': task,
'enabled': True,
}
r.hset(self.schedule_key, name, json.dumps(schedule))
def try_acquire_leadership(self):
acquired = r.setnx(self.leader_key, self.node_id)
if acquired:
r.expire(self.leader_key, self.ttl)
return True
current = r.get(self.leader_key)
if current:
current = current.decode()
if current == self.node_id:
r.expire(self.leader_key, self.ttl)
return True
return False
def get_schedules(self):
schedules = r.hgetall(self.schedule_key)
result = []
for name, data in schedules.items():
schedule = json.loads(data)
schedule['name'] = name.decode() if isinstance(name, bytes) else name
result.append(schedule)
return result
def was_executed_recently(self, schedule_name, window=60):
key = f"{self.execution_key}:{schedule_name}"
last_run = r.get(key)
if last_run:
elapsed = time.time() - float(last_run)
return elapsed < window
return False
def mark_executed(self, schedule_name):
key = f"{self.execution_key}:{schedule_name}"
r.setex(key, 3600, time.time())
def tick(self):
if not self.try_acquire_leadership():
return []
schedules = self.get_schedules()
triggered = []
for schedule in schedules:
if not schedule.get('enabled', True):
continue
if self.was_executed_recently(schedule['name'], window=60):
continue
now = datetime.now()
if self._matches_cron(schedule['cron'], now):
print(f"Leader executing: {schedule['name']}")
self.mark_executed(schedule['name'])
triggered.append(schedule)
return triggered
def _matches_cron(self, expr, dt):
parts = expr.split()
if len(parts) != 5:
return False
fields = [
(parts[0], dt.minute, 0, 59),
(parts[1], dt.hour, 0, 23),
(parts[2], dt.day, 1, 31),
(parts[3], dt.month, 1, 12),
(parts[4], dt.weekday(), 0, 6),
]
for pattern, value, lo, hi in fields:
if pattern != '*' and pattern != str(value):
if pattern.startswith('*/'):
step = int(pattern[2:])
if value % step != 0:
return False
elif '-' in pattern:
parts_range = pattern.split('-')
if not (int(parts_range[0]) <= value <= int(parts_range[1])):
return False
elif ',' in pattern:
if str(value) not in pattern.split(','):
return False
else:
return False
return True
scheduler = DistributedCronScheduler('node-1')
scheduler.register_schedule('backup', '*/1 * * * *', {'type': 'backup'})
scheduler.register_schedule('report', '*/2 * * * *', {'type': 'report'})
for _ in range(3):
jobs = scheduler.tick()
for job in jobs:
print(f" Dispatching: {job['task']}")
time.sleep(35)
Expected output:
Leader executing: backup
Dispatching: {'type': 'backup'}
Leader executing: report
Dispatching: {'type': 'report'}
Leader executing: backup
Dispatching: {'type': 'backup'}
Distributed Lock-Based Cron
#!/usr/bin/env python3
"""Per-job distributed locking for cron."""
import redis
import time
import uuid
r = redis.Redis()
class DistributedCronLock:
def __init__(self, lock_prefix='cron_job_lock', ttl=3600):
self.lock_prefix = lock_prefix
self.ttl = ttl
def execute_if_leader(self, job_name, task_func, *args):
lock_key = f"{self.lock_prefix}:{job_name}"
lock_value = str(uuid.uuid4())
acquired = r.setnx(lock_key, lock_value)
if not acquired:
print(f"[{job_name}] Another node holds the lock")
return False
r.expire(lock_key, self.ttl)
try:
print(f"[{job_name}] Acquired lock, executing")
result = task_func(*args)
return result
finally:
current = r.get(lock_key)
if current and current.decode() == lock_value:
r.delete(lock_key)
lock = DistributedCronLock()
def run_backup():
print(" Running backup...")
time.sleep(2)
return {'status': 'ok', 'files': 42}
lock.execute_if_leader('daily-backup', run_backup)
lock.execute_if_leader('daily-backup', run_backup)
Expected output:
[daily-backup] Acquired lock, executing
Running backup...
[daily-backup] Another node holds the lock
Consensus-Based Scheduling
#!/usr/bin/env python3
"""Consensus-based distributed cron using Redis sorted sets."""
import redis
import time
import json
r = redis.Redis()
class ConsensusCron:
def __init__(self, node_id, schedule_key='cron:schedule'):
self.node_id = node_id
self.schedule_key = schedule_key
self.vote_key = 'cron:votes'
def propose(self, job_name, scheduled_time):
proposal = {
'job': job_name,
'scheduled_at': scheduled_time,
'proposed_by': self.node_id,
'proposed_at': time.time(),
}
r.zadd(self.schedule_key, {json.dumps(proposal): scheduled_time})
print(f"{self.node_id} proposed {job_name} at {scheduled_time}")
def vote(self, proposal_json):
proposal = json.loads(proposal_json)
score = r.zscore(self.schedule_key, proposal_json)
if score is None:
return
vote = {
'node': self.node_id,
'job': proposal['job'],
'agree': True,
'voted_at': time.time(),
}
r.zadd(self.vote_key, {json.dumps(vote): time.time()})
def count_votes(self, proposal_json):
proposal = json.loads(proposal_json)
votes = r.zrangebyscore(
self.vote_key,
time.time() - 10,
time.time()
)
count = 0
for v in votes:
vote_data = json.loads(v)
if vote_data['job'] == proposal['job'] and vote_data['agree']:
count += 1
return count
def execute_consensus(self, job_name, scheduled_time, func):
self.propose(job_name, scheduled_time)
proposals = r.zrangebyscore(self.schedule_key, scheduled_time - 1, scheduled_time + 1)
for prop in proposals:
self.vote(prop)
for prop in proposals:
votes = self.count_votes(prop)
min_votes = 2
if votes >= min_votes:
print(f"Consensus reached for {job_name} ({votes} votes)")
func()
r.zremrangebyscore(self.schedule_key, scheduled_time - 1, scheduled_time + 1)
return
print(f"No consensus for {job_name}")
cc = ConsensusCron('node-a')
cc.propose('daily-report', time.time())
time.sleep(1)
cc.vote(json.dumps({'job': 'daily-report', 'scheduled_at': int(time.time()) - 1, 'proposed_by': 'node-b', 'proposed_at': time.time() - 1}))
Common Mistakes
1. Running Cron on Every Server
Without coordination, each server runs every job independently. Use leader election or distributed locks to ensure single execution.
2. Hardcoding Server-Specific Roles
Designating specific servers for specific jobs creates single points of failure. Use leader election so any server can take over.
3. Long Lock TTL Without Heartbeat
If the leader crashes, other servers wait for the TTL to expire before taking over. Long TTL means long downtime. Use heartbeat-based renewal.
4. Ignoring Clock Skew
Servers with different system clocks may disagree on whether a job should run. Use NTP synchronization for all servers.
5. Not Handling Network Partitions
During a network split, both partitions may elect a leader, causing duplicate execution. Use a quorum-based approach.
Practice Questions
1. What is leader election in distributed cron?
A Process where multiple servers compete to become the single scheduler. The winner runs cron jobs; losers wait. If the leader fails, another takes over.
2. How does Redis SETNX help with distributed locking?
SETNX sets a key only if it does not exist. The first server to call SETNX becomes the leader. Others see the key exists and wait.
3. What is a heartbeat in distributed cron?
A periodic signal from the leader indicating it is alive. The leader renews the lock TTL. If the heartbeat stops, the lock expires and another server becomes leader.
4. How do you prevent duplicate execution during leader failover?
Use an execution log (Redis set, database table) that records recently executed jobs. A new leader checks this log before running jobs.
Challenge
Build a distributed cron system for a 5-node cluster: leader election with Redis (10 second TTL, 5 second heartbeat), per-job execution tracking to prevent duplicates during failover, alerting when leadership changes, handling clock skew with NTP verification, and graceful degradation when Redis is unavailable (fall back to local cron).
FAQ
Mini Project: Distributed Cron Coordinator
#!/usr/bin/env python3
"""Coordinator for distributed cron with leader election."""
import redis
import time
import json
import uuid
import os
r = redis.Redis()
class DistributedCronCoordinator:
def __init__(self, node_id=None):
self.node_id = node_id or f"{os.uname().nodename}-{uuid.uuid4().hex[:4]}"
self.leader_key = 'cron:coordinator:leader'
self.jobs_key = 'cron:coordinator:jobs'
self.history_key = 'cron:coordinator:history'
self.leader_ttl = 15
def elect_leader(self):
acquired = r.setnx(self.leader_key, self.node_id)
if acquired:
r.expire(self.leader_key, self.leader_ttl)
return True
return False
def renew_leadership(self):
current = r.get(self.leader_key)
if current and current.decode() == self.node_id:
r.expire(self.leader_key, self.leader_ttl)
return True
return False
def register_job(self, name, cron_expr, job_type, data=None):
job = {
'name': name,
'cron': cron_expr,
'type': job_type,
'data': data or {},
'enabled': True,
}
r.hset(self.jobs_key, name, json.dumps(job))
def get_due_jobs(self):
jobs = r.hgetall(self.jobs_key)
due = []
for name, data in jobs.items():
name = name.decode()
job = json.loads(data)
if not job['enabled']:
continue
last_run = r.hget(self.history_key, name)
last_run = float(last_run) if last_run else 0
if time.time() - last_run >= self._interval_seconds(job['cron']):
due.append(job)
return due
def mark_run(self, job_name):
r.hset(self.history_key, job_name, time.time())
def _interval_seconds(self, cron_expr):
if cron_expr == '* * * * *':
return 60
parts = cron_expr.split()
if parts[0].startswith('*/'):
return int(parts[0][2:]) * 60
if parts[1].startswith('*/'):
return int(parts[1][2:]) * 3600
return 3600
def run(self):
if not self.elect_leader():
print(f"{self.node_id} is not the leader")
return []
self.renew_leadership()
due = self.get_due_jobs()
for job in due:
self.mark_run(job['name'])
return due
coordinator = DistributedCronCoordinator('node-1')
coordinator.register_job('health', '* * * * *', 'health_check')
coordinator.register_job('backup', '*/2 * * * *', 'backup_db')
jobs = coordinator.run()
for job in jobs:
print(f"Running {job['name']} ({job['type']})")
What's Next
Now that you understand distributed cron, explore cron in Docker containers, then learn about cron in Kubernetes.
Built by the developers of DodaTech
Doda Browser, DodaZIP & Durga Antivirus Pro