How Marin-Fray Provides a Distributed Execution Substrate
Marin-Fray delivers a distributed execution substrate by wrapping the Iris job orchestration system behind a unified Client interface, translating high-level resource and job specifications into low-level Iris protobuf messages while supporting actor hosting, coscheduling, and dynamic discovery.
The marin-community/marin repository implements Fray as a lightweight, extensible layer that turns Iris-managed clusters into a portable distributed execution substrate. By abstracting backend specifics behind a Client protocol, Fray enables users to submit jobs, host actors, and manage replicated actor groups without interacting directly with Iris RPC details.
Backend Abstraction Layer
At the core of Fray's distributed execution substrate is the FrayIrisClient class defined in lib/fray/src/fray/iris_backend.py. This class implements the Client protocol, acting as a bridge between Fray's high-level APIs and Iris's job orchestration primitives.
The abstraction allows alternative backends (such as local thread pools or custom cloud schedulers) to be swapped by implementing the same protocol, making the substrate portable across different infrastructure targets.
Resource Conversion Logic
Before submitting work to Iris, Fray must translate its native dataclasses into Iris protobuf structures. The conversion functions handle three critical mappings:
convert_resources– Transforms Fray'sResourceConfiginto IrisResourceSpecobjects, mapping device types like GPU or TPU to their Iris equivalents.convert_constraints– Converts FrayConstraintobjects into Iris placement constraints.convert_entrypoint– Packages FrayEntrypointdefinitions into IrisEntrypointstructs.
These conversion utilities ensure that high-level user specifications in fray/types.py align with the low-level requirements of the Iris backend.
Coscheduling Configuration
For multi-replica workloads requiring gang scheduling, resolve_coscheduling selects the appropriate Iris CoschedulingConfig based on device type and replica count. When running on TPUs or multi-GPU nodes, this function configures the backend to reserve resources across machines simultaneously, preventing partial allocation deadlocks.
Job Submission Pipeline
The submit method in FrayIrisClient orchestrates the transition from Fray JobRequest to Iris JobRequest. This method performs several critical operations:
- Builds the Iris protobuf request using the conversion utilities
- Applies optional wrappers for multiprocess execution or Nsight profiling
- Invokes
iris.submitto enqueue the job - Wraps the returned Iris
Jobin anIrisJobHandle
Error handling includes specific detection of duplicate job names, raising FrayJobAlreadyExists when Iris reports naming conflicts.
Job Lifecycle Management
IrisJobHandle provides a thin abstraction over Iris job state, exposing methods like status(), wait(), logs(), and terminate(). This class maps Iris's detailed state machine to Fray's simplified JobStatus enum, offering a consistent interface for monitoring distributed work regardless of backend specifics.
Actor Hosting and Group Management
Fray's substrate supports fine-grained distributed computing through an actor model layered atop Iris jobs.
Hosting Actors Within Jobs
The _host_actor function enables an actor class to run inside an Iris job replica. When invoked, it:
- Creates an
ActorServerinstance - Registers a unique endpoint following the pattern
<job_id>/<name>-<task_index> - Blocks on a shutdown event to keep the job alive
This mechanism allows driver-side services and parameter servers to coexist within the same Iris job that runs compute workers, reducing scheduling overhead.
Dynamic Actor Discovery
IrisActorGroup implements scalable actor farms by polling the Iris resolver for endpoints matching a specific prefix. This class tracks discovered actors across job replicas and provides helper methods:
discover_new– Identifies recently registered actorswait_ready– Blocks until a specified count of actors are availableshutdown– Coordinates graceful termination across the group
By leveraging Iris's endpoint registry rather than spawning separate jobs per replica, IrisActorGroup enables elastic, fault-tolerant collections of workers without overwhelming the scheduler.
Code Examples
Below are practical patterns for utilizing Fray's distributed execution substrate.
Submitting a Distributed Job
from fray import FrayIrisClient, JobRequest, ResourceConfig, GpuConfig, FrayEntrypoint
# Initialize the client against an Iris controller
client = FrayIrisClient(controller_address="iris.example.com:1234")
def train_step():
print("Executing distributed training step")
# Configure resources for 2 GPU replicas
request = JobRequest(
name="distributed_training",
entrypoint=FrayEntrypoint.from_callable(train_step),
resources=ResourceConfig(device=GpuConfig(variant="v100", count=2)),
replicas=2,
)
# Submit and monitor
handle = client.submit(request)
print(f"Job ID: {handle.job_id}")
print(f"Status: {handle.status()}")
handle.wait()
Creating an Actor Group
from fray import ActorConfig
class ParameterServer:
def __init__(self, initial_weights):
self.weights = initial_weights
def update(self, gradients):
for key, grad in gradients.items():
self.weights[key] -= 0.01 * grad
return self.weights
# Launch 4 replicas as a single Iris job with gang scheduling
ps_group = client.create_actor_group(
actor_class=ParameterServer,
name="parameter_server",
count=4,
resources=ResourceConfig(device=GpuConfig(variant="v100", count=1)),
actor_config=ActorConfig(priority=1),
initial_weights={"layer1": 0.0, "layer2": 0.0},
)
# Wait for at least one replica to be ready
ps_handle = ps_group.wait_ready(count=1)[0]
# Execute remote method call
future = ps_handle.update.remote({"layer1": 0.5, "layer2": 0.3})
updated_weights = future.result()
Hosting a Service Actor
from fray import host_actor
class Logger:
def log(self, message: str) -> None:
print(f"[LOG] {message}")
# Host within the current job context for other tasks to discover
logger = host_actor(Logger, name="experiment_logger")
# Other job replicas can resolve this endpoint via the Iris resolver
Summary
- Fray provides a distributed execution substrate by wrapping Iris behind a unified
Clientinterface, isolating users from protobuf and RPC complexities. - Resource conversion functions in
iris_backend.pymap Fray'sResourceConfigandEntrypointobjects to IrisResourceSpecand job definitions. - Coscheduling support enables gang-scheduled multi-replica jobs through the
resolve_coschedulingutility, critical for TPU and multi-GPU workloads. - Job submission occurs through
FrayIrisClient.submit, which returnsIrisJobHandleinstances for lifecycle management and status monitoring. - Actor hosting via
_host_actorruns actor classes inside job replicas, registering unique endpoints in Iris's resolver. - Dynamic discovery through
IrisActorGrouppolls the Iris resolver to build elastic, fault-tolerant actor collections without per-actor job overhead.
Frequently Asked Questions
How does Marin-Fray translate high-level job specifications into Iris-specific calls?
Fray utilizes conversion functions located in lib/fray/src/fray/iris_backend.py to transform Fray dataclasses into Iris protobuf messages. The convert_resources function maps ResourceConfig objects to Iris ResourceSpec, while convert_entrypoint packages Python callables into Iris job entrypoints. This translation layer ensures that user code remains backend-agnostic while correctly interfacing with Iris's RPC requirements.
What mechanism enables gang scheduling for multi-replica workloads?
The resolve_coscheduling function in iris_backend.py analyzes the device type and replica count to configure an Iris CoschedulingConfig. For TPU workloads or GPU gangs, this ensures all replicas are allocated simultaneously across machines, preventing partial resource acquisition that could deadlock distributed training jobs. The configuration is embedded in the JobRequest before submission to the Iris controller.
How does Fray handle actor discovery across distributed job replicas?
Fray implements dynamic actor discovery through the IrisActorGroup class, which polls the Iris resolver for endpoints matching a job-specific prefix. Rather than spawning separate Iris jobs for each actor replica, IrisActorGroup tracks endpoints formatted as <job_id>/<name>-<task_index> within a single job. This approach reduces scheduler load while enabling elastic scaling and fault tolerance through methods like discover_new and wait_ready.
Can Fray support execution backends other than Iris?
Yes. Fray defines a Client protocol that FrayIrisClient implements for the Iris backend. Alternative substrates such as local thread pools, Kubernetes operators, or custom cloud schedulers can be integrated by implementing the same protocol methods (submit, create_actor, create_actor_group, etc.). This abstraction makes the distributed execution substrate portable across different infrastructure targets without modifying user code.
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 →