Project

General

Profile

Actions

Container dispatch » History » Revision 3

« Previous | Revision 3/26 (diff) | Next »
Peter Amstutz, 12/02/2015 07:49 PM


Crunch2 dispatch

Some discussion notes from https://dev.arvados.org/issues/6429 and https://dev.arvados.org/issues/6518

Framework

Suggest writing crunch 2 job dispatcher as a new set of actors in node manger.

This would enable us to solve the question of communication between the scheduler and cloud node management (#6520).

Node manager already has a lot of the framework we will want like concurrency (can have one actor per job) and a configuration system.

Different schedulers (slurm, sge, kubernetes) can be implemented as modules similarly to how different cloud providers are supported now.

Interaction with API

More ideas:

Have a "dispatchers" table. Dispatcher processes are responsible for pinging the API server similar to how it is done for nodes to show they are alive.

A dispatcher claims a container by setting "dispatcher" field to it's UUID. This field can only be set once and that locks the record so that only the dispatcher can update it.

If a dispatcher stops pinging, the containers it has claimed should be marked as TempFail.

Dispatchers should be able to annotate containers (preferably through links) for example "I can't run this because I don't have any nodes with 40 GiB of RAM".

Retry

How do we handle failure? Is the dispatcher required to retry containers that fail, or is the dispatcher a "best effort" service and the API decides to retry by scheduling a new container?

Currently the container_uuid field only holds a single container_uuid at a time. If the API schedules a new container, does that mean any container requests associated with that container get updated with the new container?

If the container_uuid field only holds one container at a time, and container don't link back to the container requests that created, then we don't have a way to record of past attempts to fulfill this request. This means we don't have anything to check against container_count_max. A few possible solutions:

  • Make container_uuid an array of containers created to fulfill a given container request (this introduces complexity)
  • Decrement container_count_max on the request when submitting a new container
  • Compute content address of the container request and discover containers with that content address. This would conflict with "no reuse" or "impure" requests which are supposed to ignore past execution history. Could solve this by salting the content address with a timestamp; "no reuse" containers would never ever be reusable which might be fine.

I think we should distinguish between infrastructure failure and task failure by distinguishing between "TempFail" and "PermFail" in the container state. "TempFail" shouldn't count againt the container_count_max count, or alternately we only honor container_count_max for "TempFail" tasks and don't retry "PermFail".

Ideally, "TempFail" containers should retry forever, but with a backoff. One way to do the backoff is to schedule the container to run at a specific time in the future.

Scheduling

Having a field specifying "wait until time X to run this container" would be generally useful for cron-style tasks.

Updated by Peter Amstutz about 9 years ago · 26 revisions