Skip to content

Distributed Cron Scheduling — Complete Guide

DodaTech Updated 2026-06-28 8 min read

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

What happens if the Redis server goes down?

All distributed cron stops. No jobs run until Redis is restored. For high availability, use Redis Sentinel or Redis Cluster.

How quickly can a failover happen?

With a 30-second lock TTL and heartbeat every 10 seconds, failover happens within 10-30 seconds. Use shorter TTL for faster failover at the cost of more Redis calls.

Can I use Kubernetes CronJob for distributed cron?

Yes. Kubernetes CronJob runs jobs on one node in the cluster. It handles leader election internally. But it lacks some advanced scheduling features.

Should all servers run the same crontab?

Yes. If all servers run identical crontabs and use leader election, any server can become the scheduler. This simplifies configuration management.

How do I handle different timezones in distributed cron?

Use UTC for all schedules and execution tracking. Convert to local time only for logging and display. NTP ensures all servers agree on UTC.

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