Skip to content

Distributed Computing

Run computations across multiple workers with job coordination.

Enable Distributed Mode

Use reserve_jobs=True to enable job coordination:

# Single worker (default)
SessionAnalysis.populate()

# Distributed mode with job reservation
SessionAnalysis.populate(reserve_jobs=True)

How It Works

With reserve_jobs=True:

  1. Worker checks the jobs table for pending work
  2. Atomically reserves a job before processing
  3. Other workers see the job as reserved and skip it
  4. On completion, job is marked success (or error)

Multi-Process on Single Machine

# Use multiple processes
SessionAnalysis.populate(reserve_jobs=True, processes=4)

Each process:

  • Opens its own database connection
  • Reserves jobs independently
  • Processes in parallel

Multi-Machine Cluster

Run the same script on multiple machines:

# worker_script.py - run on each machine
import datajoint as dj
from my_pipeline import SessionAnalysis

# Each worker reserves and processes different jobs
SessionAnalysis.populate(
    reserve_jobs=True,
    display_progress=True,
    suppress_errors=True
)

Workers automatically coordinate through the jobs table.

Job Table

Each auto-populated table has a jobs table (~~table_name):

# View job status
SessionAnalysis.jobs

# Filter by status
SessionAnalysis.jobs.pending
SessionAnalysis.jobs.reserved
SessionAnalysis.jobs.errors
SessionAnalysis.jobs.completed

Job Statuses

Status Description
pending Queued, ready to process
reserved Being processed by a worker
success Completed successfully
error Failed with error
ignore Marked to skip

Refresh Job Queue

Sync the job queue with current key_source:

# Add new pending jobs, remove stale ones
result = SessionAnalysis.jobs.refresh()
print(f"Added: {result['added']}, Removed: {result['removed']}")

Priority Scheduling

Control processing order with priorities:

# Refresh with specific priority
SessionAnalysis.jobs.refresh(priority=1)  # Lower = more urgent

# Process only high-priority jobs
SessionAnalysis.populate(reserve_jobs=True, priority=3)

Error Recovery

Handle failed jobs:

# View errors
errors = SessionAnalysis.jobs.errors
for job in errors.to_dicts():
    print(f"Key: {job}, Error: {job['error_message']}")

# Clear errors to retry
errors.delete()
SessionAnalysis.populate(reserve_jobs=True)

Orphan Detection

A worker killed outright โ€” SIGKILL, an OOM kill, a lost node โ€” leaves its job reserved with no error row, and no other worker will take it. Recovering it takes an explicit refresh:

# Refresh with orphan timeout (seconds)
SessionAnalysis.jobs.refresh(orphan_timeout=3600)

Reserved jobs older than the timeout are reset to pending. Nothing does this on your behalf: orphan_timeout has no default and no configuration key, and the refresh that populate() runs for itself never passes one. Schedule the call, or run it when a worker is known to have died.

Choose the timeout above the longest make() you expect. Jobs are reclaimed by the age of their reservation, not by checking whether the worker is alive, so a computation that outlives the timeout is re-queued while it is still running and a second worker can start the same key.

A worker that receives SIGTERM is a different case โ€” it records an error and needs no orphan handling. See Signals and Interruption.

Configuration

import datajoint as dj

# Auto-refresh on populate (default: True)
dj.config.jobs.auto_refresh = True

# Keep completed job records (default: False)
dj.config.jobs.keep_completed = True

# Stale job timeout in seconds (default: 3600)
dj.config.jobs.stale_timeout = 3600

# Default job priority (default: 5)
dj.config.jobs.default_priority = 5

# Track code version (default: None)
dj.config.jobs.version_method = "git"

Populate Options

Option Default Description
reserve_jobs False Enable job coordination
processes 1 Number of worker processes
max_calls None Limit jobs per run
display_progress False Show progress bar
suppress_errors False Continue on errors
priority None Filter by priority
refresh None Force refresh before run

Example: Cluster Setup

# config.py - shared configuration
import datajoint as dj

dj.config.jobs.auto_refresh = True
dj.config.jobs.keep_completed = True
dj.config.jobs.version_method = "git"

# worker.py - run on each node
from config import *
from my_pipeline import SessionAnalysis

while True:
    result = SessionAnalysis.populate(
        reserve_jobs=True,
        max_calls=100,
        suppress_errors=True,
        display_progress=True
    )
    if result['success_count'] == 0:
        break  # No more work

See Also