Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

tfmesos2

tfmesos2 runs TensorFlow tasks as an Apache Mesos framework. A Python process registers the framework, accepts matching resource offers, and starts one container per task. The containers report their port and readiness to the local API; TensorFlow can then use the generated cluster definition.

Goals

  • run parameter-server and worker tasks in parallel on Mesos
  • request CPU, memory, and optional GPU resources per job
  • use Docker containers as the task runtime
  • keep local tests reproducible without a live cluster

This documentation describes the behavior of the current Python implementation. Production values such as credentials, images, and hostnames belong in environment variables or a secret-management system.

Architecture

Execution flow

  1. tfmesos2.cluster(...) normalizes job dictionaries into Job objects.
  2. TensorflowMesos creates task IDs and registers a Mesos client.
  3. For each offer, CPU, memory, GPUs, and MESOS_ATTRIBUTES are checked.
  4. A Mesos task starts python3 -m tfmesos2.server <task-id> <client-ip>.
  5. The task reports a random local port and waits for cluster metadata.
  6. The API returns cluster_def; the task starts tf.distribute.Server.
  7. Once all tasks are ready, the framework is suppressed and accepts no further offers.

Component boundaries

  • Scheduler process: owns the Mesos framework, task queue, resource matching, and readiness state.
  • Mesos master and agents: place and run the TensorFlow containers.
  • Task callback API: exchanges task ports, readiness, and cluster metadata.
  • TensorFlow task: starts the gRPC server and executes the user workload.

Failure boundaries

Terminal Mesos states are never treated as successful readiness. wait_until_ready() aborts on TASK_FAILED, TASK_LOST, or TASK_KILLED and has a configurable timeout. Status updates for unknown task IDs are ignored and logged.

Configuration

Mesos

export MESOS_MASTER=mesos-master.example.test:5050
export MESOS_SSL=true
export MESOS_VERIFY_SSL=true
export MESOS_USERNAME=principal-from-secret-store
export MESOS_PASSWORD=secret-from-secret-store
export MESOS_FRAMEWORK_ROLE=tensorflow

MESOS_SSL=true uses HTTPS for Mesos URLs; otherwise HTTP is used. MESOS_VERIFY_SSL controls certificate verification and defaults to enabled. Disable verification only in a controlled internal test environment.

Resources and image

Jobs must contain at least name and num; cpus, mem, and gpus are optional. The Docker image is read from DOCKER_IMAGE, defaulting to avhost/tensorflow-mesos:latest. Offer placement can be constrained with MESOS_ATTRIBUTES in the form key:value.

Task API

The task uses TFMESOS2_CLIENT_USERNAME, TFMESOS2_CLIENT_PASSWORD, TFMESOS2_VERIFY_SSL, and TFMESOS2_REQUEST_TIMEOUT. The default HTTP request timeout is 30 seconds.

Operations and Monitoring

Local verification

python3 -m compileall -q tfmesos2
python3 -m unittest discover -s tests -v
make -C docs build

Deployment

The documentation can be published to the gh-pages branch through the repository’s deployment worktree:

make -C docs deploy

The target builds the book first, stages the generated site in a temporary worktree, and pushes only when the rendered output changed. The defaults are origin as the remote and gh-pages as the branch. Override them when needed:

make -C docs deploy DEPLOY_REMOTE=origin DEPLOY_BRANCH=gh-pages

The deployment worktree must not already exist. The command requires Git push permission for the configured remote.

A real cluster requires the task image to contain tfmesos2.server and TensorFlow. The Mesos master must be reachable from the task network, and the callback address (client_ip:port) must be reachable from the started containers.

Observability

Mesos task states and declined offers are written to the Python log. The API provides /v0/status, /v0/task/<task-id>/job, /v0/task/<task-id>/port/<port>, and /v0/task/<task-id>. After successful initialization, the framework is suppressed.

Resources

Choose CPU and memory values with enough headroom for TensorFlow, gRPC, and the container runtime. GPU offers are handled as SET or SCALAR resources; the Docker parameters depend on the configured GPU vendor.

HTTP API

The API listens on port 11000 by default and is called without TLS inside the task network. Client calls send the configured HTTP Basic Auth values; however, the current Flask API does not enforce authentication itself. In production, keep it behind an appropriately protected internal network boundary.

MethodPathPurpose
GET/v0/statusReturns ok once all active tasks are TASK_RUNNING
GET/v0/task/<id>/jobReturns the job name, task index, and current cluster_def
PUT/v0/task/<id>/port/<port>Reports a task port; the port must be 1–65535
PUT/v0/task/<id>Marks a task as initialized
GET/v0/download/<filename>Serves a file from the working directory

Unknown task IDs return 404 for mutation requests; invalid ports return 400. The download route rejects paths containing ...

Troubleshooting

Readiness timeout

First check the Mesos task states and callback-address reachability. A task can be TASK_RUNNING while its HTTP callback is still failing. Also verify that the container reports the correct port.

No matching offers

Compare the CPU, memory, and GPU requirements with the offer. With MESOS_ATTRIBUTES, the offer text must match the key:value constraint. Invalid values are rejected instead of being silently interpreted.

TLS or authentication errors

Make sure MESOS_SSL matches the listener, credentials are complete, and the CA certificate is available to the process. MESOS_VERIFY_SSL=false is not a substitute for a correct CA configuration.

TensorFlow does not start

Verify that TensorFlow can be imported inside the task image and that cluster_def contains every expected job with reachable host:port entries. The scheduler can be tested without TensorFlow; runtime validation requires a compatible image.