| name | labtasker |
| description | Use when helping users distribute ML experiment tasks, convert for-loop scripts into task queues, manage parallel GPU workers, handle failures/retries, or query/filter/update/delete experiment results with Labtasker |
Labtasker
Overview
Labtasker replaces for loops in ML experiment scripts with a task queue, enabling parallelization, failure handling, and task management with minimal code changes.
Links: GitHub | Docs
Core Concepts
- Queue: Named task bucket; all tasks belong to one queue; workers pull from it.
- Task fields:
task_id, task_name, status (pending/running/success/failed/cancelled), args (dict), metadata (dict), summary (dict, always a dict — never None — even if no metrics written), priority, max_retries, retries, created_at, last_modified, start_time (None until task starts running), worker_id (None until a worker claims the task).
- Loop: Continuously fetches pending tasks in descending priority order (highest
priority value first); exits automatically when no more pending tasks match.
- Retries:
retries starts at 0 on the first attempt. The loop auto-increments it on each task failure. When retries reaches max_retries, the task is permanently marked failed. Use reset_pending=True to reset both status and retries when re-queuing.
- No More, No Less: Submitted
args keys must EXACTLY match what the run script/function declares as required. Extra = inconsistent records. Missing = fetch failure. Exception: with pass_args_dict=True + required_fields, only the listed keys are required — extra keys in args are accessible via args.get(key, default) for variable-arg tasks.
Setup
pip install labtasker
labtasker-server serve &
labtasker init
labtasker queue create-from-config
1. Submitting Tasks
CLI
labtasker task submit -- --lr=0.001 --model=resnet
labtasker task submit \
--name grid_search \
--metadata '{"tags": ["v1", "experimental"], "note": "baseline run"}' \
--priority 10 \
--max-retries 5 \
-- --lr=0.001 --model=resnet
for lr in 0.001 0.01 0.1; do
for model in resnet vit; do
labtasker task submit \
--name grid_search \
--metadata '{"tags": ["v1"]}' \
-- --lr=$lr --model=$model
done
done
labtasker task submit --args '{"arg1": 0, "arg2": 3}'
labtasker task submit -- --optimizer.lr=0.01 --optimizer.momentum=0.9
Special characters in args — always use --key=value form:
labtasker task submit -- --value=-1 --label="" --desc="hello world"
Python
import labtasker
for lr in [0.001, 0.01, 0.1]:
for model in ["resnet", "vit"]:
labtasker.submit_task(
task_name="grid_search",
args={"arg1": lr, "arg2": model},
metadata={
"tags": ["v1", "experimental"],
"note": "baseline run",
},
max_retries=3,
priority=0,
)
labtasker.submit_task(
args={
"args_a": {"a": 1, "b": "boy"},
"args_b": {"foo": 2, "bar": "baz"},
},
metadata={"tags": ["sweep-v2"]},
)
2. Running Tasks (Loop)
CLI
labtasker loop -- python train.py --lr '%(lr)' --model '%(model)'
labtasker loop \
--extra-filter '"v1" in list(metadata.tags)' \
-- python train.py --lr '%(lr)' --model '%(model)'
CUDA_VISIBLE_DEVICES=0 labtasker loop -- python train.py --lr '%(lr)'
CUDA_VISIBLE_DEVICES=0 labtasker loop -- python train.py --lr '%(lr)' --model '%(model)' &
CUDA_VISIBLE_DEVICES=1 labtasker loop -- python train.py --lr '%(lr)' --model '%(model)' &
CUDA_VISIBLE_DEVICES=2 labtasker loop -- python train.py --lr '%(lr)' --model '%(model)' &
CUDA_VISIBLE_DEVICES=3 labtasker loop -- python train.py --lr '%(lr)' --model '%(model)' &
wait
LABTASKER_TASK_SCRIPT=$(mktemp)
cat <<'EOF' > "$LABTASKER_TASK_SCRIPT"
LOG_DIR=/logs/%(dataset)/%(model)
python train.py --dataset %(dataset) --model %(model) --log-dir $LOG_DIR
EOF
labtasker loop --script-path $LABTASKER_TASK_SCRIPT
%(key) accesses top-level args keys only. If args has a nested dict (e.g., {"optimizer": {"lr": 0.01}}), use the Python loop instead — CLI cannot interpolate nested values.
Python
import labtasker
from labtasker import Required
@labtasker.loop()
def main(arg1: int = Required(), arg2: int = Required()):
result = arg1 + arg2
print(f"The result is {result}")
if __name__ == "__main__":
main()
@labtasker.loop(
extra_filter='"v1" in list(metadata.tags)'
)
def main(lr: float = Required(), model: str = Required()):
train(lr=lr, model=model)
@labtasker.loop(required_fields=["lr", "model"], pass_args_dict=True)
def main(args):
train(lr=args["lr"], model=args["model"])
from typing import Any, Dict
from typing_extensions import Annotated
from dataclasses import dataclass
@dataclass
class ArgsGroupA:
a: int
b: str
@dataclass
class ArgsGroupB:
foo: int
bar: str
@labtasker.loop()
def main(
args_a: Annotated[Dict[str, Any], Required(resolver=lambda a: ArgsGroupA(**a))],
args_b=Required(resolver=lambda b: ArgsGroupB(**b)),
):
print(f"got args_a: {args_a}")
print(f"got args_b: {args_b}")
@labtasker.loop()
def main(optimizer=Required(resolver=lambda x: x)):
lr = optimizer["lr"]
momentum = optimizer.get("momentum", 0.9)
The loop exits automatically when no more pending tasks match the filter.
Dot-separated keys are nested: --foo.bar=1 in CLI becomes args["foo"]["bar"] = 1 in Python. Access nested args correctly: task_info().args["foo"]["bar"], not task_info().args["foo.bar"].
3. Querying (Listing) Tasks
CLI
labtasker task ls -s pending
labtasker task ls -s running
labtasker task ls -s failed
labtasker task ls -s success
labtasker task ls -f 'args.lr > 0.01' -q --no-pager
labtasker task ls -f 'task_name == "grid_search"' -q --no-pager
labtasker task ls -f '"v1" in list(metadata.tags)' -q --no-pager
labtasker task ls -f 'created_at >= date("3 hours ago")' -q --no-pager
labtasker task ls -f 'args.lr > 0.01 and "v1" in list(metadata.tags)' -q --no-pager
labtasker task ls -f 'regex(task_name, "^grid-.*")' -q --no-pager
labtasker task ls -s success -f 'summary.acc > 0.9' --no-pager
labtasker task ls -f 'status == "failed" or status == "pending"' --no-pager
labtasker task ls -s success -S 'created_at:desc' --no-pager
labtasker task ls -s success -S 'summary.acc:desc' --no-pager
Filter operators — Python syntax: ==, >, <, >=, <=, in, and, or, regex(), date()
NOT supported: !=, not, not in (three-valued logic). Workaround for status: use -s STATUS flag. For string negation: use regex().
-s and -f are ANDed: using both narrows results to tasks matching both conditions.
Default --limit is 100. For batch operations (cancel/delete/update all matching tasks), always pass --limit 10000 or a large value to avoid silently missing tasks beyond 100.
Python
import labtasker
response = labtasker.ls_tasks(status="failed")
tasks = response.content
response = labtasker.ls_tasks(
extra_filter='args.lr > 0.01 and "v1" in list(metadata.tags)',
limit=100,
)
response = labtasker.ls_tasks(
status="success",
extra_filter='summary.acc > 0.9',
limit=1000,
)
tasks_sorted = sorted(
response.content,
key=lambda t: t.summary.get("acc") or 0,
reverse=True,
)
for task in response.content:
print(task.task_id, task.status, task.args, task.metadata, task.summary)
print(task.created_at, task.worker_id, task.retries)
4. Task Summary (Write and Read)
task.summary is a dict stored per-task. Jobs write metrics into it; you read it back after tasks complete.
Writing summary inside the loop
Use labtasker.report_task_status() to persist metrics. This call is optional — if you simply return from the function without calling it, the loop auto-marks the task success with no summary. Call it when you want to record metrics.
import labtasker
from labtasker import Required
@labtasker.loop()
def main(lr: float = Required(), model: str = Required()):
acc, loss = train(lr=lr, model=model)
labtasker.report_task_status(
task_id=labtasker.task_info().task_id,
status="running",
summary={"acc": acc, "loss": loss},
)
To finish a task early with a final status from deep inside your code — use labtasker.finish(). It stops the heartbeat, writes summary.json to the local log directory, and reports to the server. The loop will not overwrite a status already set by finish(). Only accepts "success" or "failed".
Retry behavior: finish(status="failed") behaves like an exception — the server increments retries and resets the task to pending if retries < max_retries (so it WILL be retried). To permanently stop retries without marking as failed, use report_task_status(status="cancelled") instead.
Idempotent: calling finish() twice is safe — the second call is silently skipped.
@labtasker.loop()
def main(lr: float = Required(), model: str = Required()):
acc = train(lr=lr, model=model)
if acc < 0.01:
labtasker.finish(status="failed", summary={"acc": acc, "reason": "too_low"})
return
labtasker.finish(status="success", summary={"acc": acc})
To cancel a task mid-run (e.g., on NaN loss) — use report_task_status(status="cancelled") then return. The loop preserves the cancelled status and will NOT retry it:
@labtasker.loop()
def main(lr: float = Required(), model: str = Required()):
task = labtasker.task_info()
if task.retries == 0:
data = expensive_preprocess()
else:
data = load_cached_data()
loss = train_step(data)
if math.isnan(loss):
labtasker.report_task_status(
task_id=task.task_id,
status="cancelled",
)
return
Reading summaries after tasks complete
import labtasker
response = labtasker.ls_tasks(status="success", limit=1000)
for task in response.content:
print(
f"task_id={task.task_id} "
f"args={task.args} "
f"summary={task.summary}"
)
results = [
{"lr": t.args["lr"], "model": t.args["model"], "acc": t.summary.get("acc")}
for t in response.content
]
results.sort(key=lambda x: x["acc"] or 0, reverse=True)
for r in results:
print(r)
5. Updating Tasks
Update semantics
CLI -u 'field=value' patches individual dot-nested fields: -u 'args.lr=0.005' sets only args.lr, all other args untouched.
Python TaskUpdateRequest merges by default: only the keys you include inside a dict field are written; other existing keys inside that dict are preserved. Unspecified top-level optional fields (args, metadata, summary, priority, etc.) are not touched. To replace a dict field entirely (overwrite, not merge), add it to replace_fields: TaskUpdateRequest(**{"_id": id, "args": {...}, "replace_fields": ["args"]}).
reset_pending=True sets status → pending AND resets retries → 0. Setting "status": "pending" in TaskUpdateRequest alone changes status but does not reset the retry counter. Always prefer reset_pending=True when re-queuing failed tasks.
CLI
labtasker task update --id <task_id> -u 'args.lr=0.005'
labtasker task update --id <task_id> -u 'args.lr=0.005' -u 'metadata.tags=["v2"]'
labtasker task ls -f 'status == "failed" and args.model == "resnet"' -q --limit 10000 \
| xargs -I{} labtasker task update --id {} -u 'args.lr=0.005' --reset-pending --quiet
labtasker task update -s failed -f 'args.model == "resnet"' -u 'args.lr=0.005' --reset-pending --quiet
labtasker task ls -s pending -f 'args.lr > 0.05' -q --limit 10000 \
| xargs -I{} labtasker task update --id {} -u 'status=cancelled' --quiet
labtasker task ls -f '"urgent" in list(metadata.tags)' -s pending -q --limit 10000 \
| xargs -I{} labtasker task update --id {} -u 'priority=100' --quiet
Python
import labtasker
from labtasker.api_models import TaskUpdateRequest
labtasker.update_tasks([
TaskUpdateRequest(**{"_id": task_id, "args": {"lr": 0.005}})
])
response = labtasker.ls_tasks(status="failed")
updates = [
TaskUpdateRequest(**{
"_id": task.task_id,
"args": {"lr": 0.005},
})
for task in response.content
]
if updates:
labtasker.update_tasks(updates, reset_pending=True)
response = labtasker.ls_tasks(
status="pending",
extra_filter='args.lr > 0.05',
)
updates = [
TaskUpdateRequest(**{"_id": task.task_id, "status": "cancelled"})
for task in response.content
]
if updates:
labtasker.update_tasks(updates)
@labtasker.loop()
def main(lr: float = Required(), model: str = Required()):
train(lr=lr, model=model)
task = labtasker.task_info()
labtasker.update_tasks([
TaskUpdateRequest(**{"_id": task.task_id, "metadata": {"run_note": "done"}})
])
6. Deleting Tasks
CLI
labtasker task delete <task_id>
labtasker task ls -f 'created_at < date("10 minutes ago")' -q --limit 10000 | labtasker task delete -y
labtasker task ls -s failed -q --limit 10000 | labtasker task delete -y
labtasker task ls -f 'regex(task_name, "^debug-.*")' -q --limit 10000 | labtasker task delete -y
Python
import labtasker
labtasker.delete_task(task_id="<task_id>")
response = labtasker.ls_tasks(
extra_filter='created_at < date("10 minutes ago")'
)
for task in response.content:
labtasker.delete_task(task_id=task.task_id)
Quick Reference
| Operation | CLI | Python |
|---|
| Submit | labtasker task submit -- --k=v | labtasker.submit_task(args={...}, metadata={...}) |
| Submit with tags | --metadata '{"tags": ["v1"]}' | metadata={"tags": ["v1"]} |
| Run loop | labtasker loop -- cmd '%(arg)' | @labtasker.loop() + Required() |
| Run with filter | --extra-filter 'filter' | @labtasker.loop(extra_filter='...') |
| Parallel workers | CUDA_VISIBLE_DEVICES=N labtasker loop ... & | run multiple processes |
| List by status | labtasker task ls -s failed | labtasker.ls_tasks(status="failed") |
| List with filter | labtasker task ls -f 'expr' | labtasker.ls_tasks(extra_filter='expr') |
| Sort | labtasker task ls -S created_at:desc | sorted(response.content, key=lambda t: ...) |
| Write intermediate summary | — | labtasker.report_task_status(task_id=..., status="running", summary={...}) |
| Finish task (final status) | — | labtasker.finish(status="success"|"failed", summary={...}) |
| Update | labtasker task update --id X -u 'args.k=v' | labtasker.update_tasks([TaskUpdateRequest(...)]) |
| Cancel | ... -q | xargs ... -u 'status=cancelled' | TaskUpdateRequest(**{"_id": id, "status": "cancelled"}) |
| Retry failed | --reset-pending | update_tasks(updates, reset_pending=True) |
| Delete | labtasker task ls -s failed -q | labtasker task delete -y | labtasker.delete_task(task_id=...) |
| Current task | — | labtasker.task_info() |
API Signatures
For queue and worker management, see the full documentation.
CLI
labtasker task submit
[ARGS...] # task args as CLI flags after -- e.g. -- --lr=0.01 --model=vit
[--args JSON_STR] # alternative: pass args as JSON dict string
[--name NAME] # task name for identification
[--metadata JSON_STR] # e.g. '{"tags": ["v1"], "note": "baseline"}'
[--max-retries INT] # retry attempts on failure (default: 3)
[--priority INT] # higher = higher priority (default: 0)
labtasker loop
[-- CMD ARGS] # command with %(key) placeholders; %(key) = top-level args only
[--script-path FILE] # path to script file with %(key) placeholders (for multi-line logic)
[-f / --extra-filter EXPR] # Python query string to filter which tasks to run
# For nested dict args, use the Python @labtasker.loop() instead of CLI loop
labtasker task ls
[--id TASK_ID] # filter by task ID
[--name TASK_NAME] # filter by task name
[-s / --status STATUS] # pending|running|success|failed|cancelled
[-f / --extra-filter EXPR] # Python query string (supports args.field, metadata.field, summary.field)
[-q / --quiet] # print task IDs only (for piping)
[--no-pager] # disable pager output
[--limit INT] # max results (default: 100)
[-S / --sort FIELD:asc|desc] # e.g. -S created_at:desc or -S summary.acc:desc
labtasker task update
[--id TASK_ID] # filter by task ID
[--name TASK_NAME] # filter by task name
[-s / --status STATUS] # filter by status
[-f / --extra-filter EXPR] # filter by Python query
[-u / --update 'field=value'] # patch a field (repeatable); dot-notation for nested: -u args.lr=0.01
[-- field=value ...] # positional update syntax (alternative to -u)
[--reset-pending] # set status=pending AND reset retries=0
[-q / --quiet] # skip confirmations (for scripts/pipes)
labtasker task delete
[TASK_IDS...] # task IDs to delete (positional, or piped from stdin)
[-y / --yes] # skip confirmation prompt
Python
labtasker.submit_task(
task_name: str = None,
args: Dict = None,
metadata: Dict = None,
max_retries: int = 3,
priority: int = 0,
) -> TaskSubmitResponse
@labtasker.loop(
extra_filter: str = None,
required_fields: List[str] = None,
pass_args_dict: bool = False,
)
def main(
param: Type = Required(),
param: Annotated[T, Required(resolver=fn)] = ...,
param=Required(resolver=fn),
):
...
labtasker.task_info() -> Task
labtasker.finish(
status: str,
summary: Dict = None,
skip_if_no_labtasker: bool = True,
) -> None
labtasker.report_task_status(
task_id: str,
status: str,
summary: Dict = None,
) -> None
labtasker.ls_tasks(
task_id: str = None,
task_name: str = None,
status: str = None,
extra_filter: str = None,
limit: int = 100,
) -> TaskLsResponse
labtasker.update_tasks(
task_updates: List[TaskUpdateRequest],
reset_pending: bool = False,
) -> TaskLsResponse
TaskUpdateRequest(**{
"_id": str,
"status": str,
"args": Dict,
"metadata": Dict,
"summary": Dict,
"priority": int,
"max_retries": int,
"replace_fields": List[str],
})
labtasker.delete_task(task_id: str) -> None
labtasker.update_tasks([
TaskUpdateRequest(**{"_id": task_id, "status": "cancelled"})
])
Filter Query Syntax
# Comparison
field == value field > value field >= value
field < value field <= value
# Membership (array field)
"val" in list(field)
# Logic
expr and expr expr or expr
# Functions
date("3 hours ago") # datetime: "X ago" means X before now
regex(field, "^pattern.*") # regex match on string field
# Dot-notation works on any field:
args.lr > 0.01
metadata.group == "A"
summary.acc > 0.9
created_at >= date("7 days ago")
# NOT supported: != not not in
# Workaround for status negation: use -s STATUS flag (CLI) or status= param (Python)
# Workaround for OR over statuses: use status == "pending" or status == "failed"
Event system (SSE): For real-time task lifecycle events and workflow automation,
see Advanced Features.
Workflow
- Identify the serial loop in the user's script; choose CLI or Python based on their code style.
- Submit script: iterate parameter grid, call
submit_task/labtasker task submit with matching args and any metadata tags needed.
- Run script: declare exactly the same keys as
Required() params or %(key) placeholders — no more, no less.
- For parallel runs: each worker independently executes the run script; all pull from the same queue.
- Env vars must go outside the
labtasker loop command (CLI).
- To retry failures: use
update_tasks(updates, reset_pending=True) (resets status AND retries); or --reset-pending in CLI.
- To review results:
ls_tasks(status="success") and inspect .summary and .args on each task.