How Marin-Iris Manages Distributed Job Execution: Architecture and Implementation
Marin-Iris orchestrates distributed training and inference workloads across local GPU, Cloud GPU/TPU, and federated Kubernetes clusters through a controller-database-scheduler architecture that guarantees state consistency by preventing backends from writing directly to the job database.
Marin-Iris serves as the job-orchestration layer within the marin-community/marin repository, driving distributed execution through a deterministic scheduling pipeline. Understanding how Marin-Iris manages distributed job execution reveals a design that separates state management from compute provisioning, ensuring reliable coordination across heterogeneous clusters.
Architecture Overview
Marin-Iris consists of three tightly-coupled components that coordinate distributed workloads:
Controller ▶ Database
Persists job metadata, configuration, and task state in a SQL database. All scheduling decisions derive from a consistent snapshot of this database, with schema definitions for jobs_table and job_workdir_files_table located in iris/cluster/schema.py.
Scheduler ▶ Decision Engine
Runs a pure-Python scheduling pipeline in run_scheduling_decision within iris/cluster/controller/backend.py (lines 47-65). This engine expands database snapshots into placement decisions through availability enrichment, gating, ordering, assignment, and pre-emption phases.
Backend ▶ Workers
Implements the TaskBackend protocol defined in iris/cluster/controller/backend.py (lines 67-74). Backends such as native Docker, Kubernetes, Kueue, or Slurm expose capacity and execute tasks via a thin RPC layer, but never write directly to the database.
Job Submission Flow
CLI and Client Requests
Users initiate jobs via the command line:
uv run iris --cluster=marin job run \
--cpu=4 --memory=8G \
--extra=cpu \
--region=us-east5 \
python train.py
The CLI marshals a controller_pb2.Controller.LaunchJobRequest and transmits it to the controller, as documented in docs/tutorials/train-an-lm.md.
Database Insertion and Configuration
The submit function calls insert_job_and_config in iris/cluster/controller/ops/job.py (lines 57-73 and 84-92). This transaction writes the jobs row and associated configuration, resolves the priority band, and records the authenticated submitter for root jobs.
Task Materialization
Following successful insertion, the controller invokes _materialize_tasks in iris/cluster/controller/ops/job.py (lines 64-85). This creates a contiguous block of PENDING task rows (one per replica) by reserving a priority-insertion base and bulk-inserting tasks.
Audit Logging
The insertion registers an audit event (job_submitted) emitted once the transaction commits (lines 48-53).
The Deterministic Scheduling Pipeline
Every controller tick builds an in-memory snapshot of the current database state and passes it to the scheduler. The pipeline executes entirely without I/O:
-
Capacity Enrichment:
run_scheduling_decisionenriches worker attributes withavailability:<variant>markers if the backend supplies per-zone accelerator capability information (lines 68-78). -
Gating:
apply_scheduling_gatesfilters tasks blocked by resource caps, quotas, or expiration (lines 84-88). -
Ordering:
compute_scheduling_orderproduces a deterministic priority sequence for remaining tasks (lines 108-109). -
Assignments:
scheduler.find_assignmentsexecutes the core placement algorithm, matching tasks to workers based on constraints, priority band, and capacity (lines 60-62). Theapply_placementsfunction collects results (lines 48-64). -
Pre-emption:
apply_preemptionsgenerates pre-empt plans for higher-priority jobs requiring slots, committed atomically with assignments (lines 110-120). -
Result: Returns a
ScheduleResultcontaining assignments, preemptions, unschedulable tasks, diagnostics, and enriched context (lines 19-36).
Backend Interaction Protocol
The controller communicates with compute providers through the TaskBackend protocol:
-
observe: Publishes provider status including worker liveness, capacity, and pending hints as aBackendObservation(lines 91-93). -
schedule: Invokes the pure decision pipeline for backends performing their own placement (e.g., Kueue). For backends allowing Iris to schedule tasks directly, this returns empty and the controller drives placement. -
runtime_image: Resolves container images with backend-specific fallback defaults (lines 82-88).
This architecture guarantees state consistency because all mutations funnel through controller transactions; backends remain read-only regarding database state.
Federated Execution and Hierarchical Jobs
Iris supports federated clusters where parent jobs launch child jobs on remote clusters. The cluster argument in insert_job_and_config routes to the child cluster's backend while the parent cluster materializes tasks locally. The controller mirrors child state back into the parent database, enabling unified visibility across clusters (lines 68-71).
Priority Bands and Resource Management
Jobs specify a PRIORITY_BAND (interactive, batch, or INHERIT). The resolve_priority_band function normalizes this at ingestion in iris/cluster/controller/ops/job.py (lines 88-99), ensuring stored rows contain concrete bands. The scheduler uses these bands to order tasks and enforce band-specific capacity limits.
Practical Code Examples
Submitting via CLI
uv run iris --cluster=marin job run \
--cpu=4 --memory=8G \
--extra=cpu \
--region=us-east5 \
python train.py
This creates a LaunchJobRequest handled by the submit path in iris/cluster/controller/ops/job.py.
Programmatic Submission
from iris.client import IrisClient
from iris.rpc import controller_pb2
client = IrisClient(cluster="marin")
req = controller_pb2.Controller.LaunchJobRequest(
job_name="example",
resources=controller_pb2.ResourceSpecProto(
cpu_millicores=4000,
memory_bytes=8*1024**3
),
command=["python", "train.py"],
)
job_id = client.submit_job(req) # Calls submit → insert_job_and_config
print(f"Submitted job {job_id}")
Inspecting Job Status
info = client.get_job_info(job_id)
print(f"State: {info.state}") # TASK_STATE_PENDING, RUNNING, etc.
print(f"Tasks: {info.num_tasks}")
Summary
- Marin-Iris separates concerns into a database-backed controller, pure-Python scheduler, and pluggable backends to manage distributed job execution.
- Job submission flows through
insert_job_and_configand_materialize_tasksiniris/cluster/controller/ops/job.py, creating atomic database transactions with audit logging. - The scheduling pipeline executes six deterministic phases—enrichment, gating, ordering, assignment, pre-emption, and result compilation—without database I/O via
run_scheduling_decision. - The
TaskBackendprotocol abstracts Docker, Kubernetes, Kueue, and Slurm while enforcing a strict boundary preventing backends from writing to the controller database. - Federated execution enables hierarchical jobs across clusters with unified state mirroring in the parent database.
- Priority bands (interactive, batch, inherit) resolve at ingestion to enforce deterministic scheduling order and capacity limits.
Frequently Asked Questions
How does Marin-Iris ensure database consistency across distributed workers?
The controller maintains exclusive write access to the SQL database. Backends implement the TaskBackend protocol and communicate via RPC, but they never write directly to the database. All state changes flow through the controller's transaction logic in iris/cluster/controller/ops/job.py, ensuring a single source of truth for job metadata and task states.
What scheduling algorithm does Marin-Iris use for task placement?
Marin-Iris uses a multi-phase deterministic algorithm implemented in run_scheduling_decision within iris/cluster/controller/backend.py. The pipeline enriches capacity data, applies scheduling gates for resource constraints, computes priority ordering, finds assignments matching tasks to workers, and applies preemptions when higher-priority jobs require resources.
Can Marin-Iris schedule jobs across multiple Kubernetes clusters?
Yes. Marin-Iris supports federated execution where parent jobs launch child jobs on remote clusters using the cluster argument in insert_job_and_config. The controller mirrors child cluster state back into the parent database, providing a unified view of tasks across local and remote infrastructure.
How are job priorities handled in the Marin-Iris scheduler?
Jobs specify a PRIORITY_BAND (interactive, batch, or INHERIT) which resolve_priority_band normalizes to a concrete band during ingestion in iris/cluster/controller/ops/job.py (lines 88-99). The scheduler uses these bands to determine task ordering and enforce band-specific capacity limits during the assignment phase.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →