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:
- Worker checks the jobs table for pending work
- Atomically reserves a job before processing
- Other workers see the job as reserved and skip it
- 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¶
- Run Computations โ Basic populate usage
- Handle Errors โ Error recovery patterns
- Monitor Progress โ Tracking job status