Skip to content

About

Consistent Rich progress bars, smoothed ETA and per-worker status for any multiprocessing or asyncio batch job

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Repository files navigation

typer-rich-progress

A small library that gives any multiprocessing or asyncio batch job a consistent Rich progress display — overall bar, smoothed ETA, throughput, and a per-worker status table — in a few lines.

Reference documentation and background articles: batch-processing.com.

The problem

Every batch CLI ends up growing its own progress code, and it is nearly always wrong in the same three ways.

The ETA is computed as elapsed / completed extrapolated over what is left. That is an average over the whole run, so a job whose per-item cost varies — a directory of rasters where the first hundred tiles are 2 MB and the next hundred are 2 GB — reports a confidently wrong estimate for its entire second half, and only becomes honest once it is nearly finished.

The display has no idea what the workers are actually doing. A bar that says 412/5000 tells you nothing about whether one worker has been stuck on the same 4 GB mosaic for ten minutes.

And it animates unconditionally. Redirect the output to a log, or run it in CI, and you get several megabytes of cursor-movement escape sequences wrapped around the two lines you wanted.

typer-rich-progress is the small piece that solves those three things once.

Install

Not published to an index — clone it:

git clone https://github.com/batch-processing-geospatial-cli-tools/typer-rich-progress.git
cd typer-rich-progress
uv sync --all-extras
uv run pytest

Or with a plain virtualenv:

python -m venv .venv
.venv/bin/pip install -e ".[typer]"

The core needs only rich. Typer is an optional extra.

Usage

Track a loop

from typer_rich_progress import BatchProgress

with BatchProgress(total=len(paths), description="Warping") as progress:
    for path in progress.track(paths):
        warp(path)
Warping ━━━━━━━━━━━━━━━━╺━━━━━━━ 412/620  66.5%  1m14s elapsed | ETA 38.2s | 5.53/s

Run it across processes

from typer_rich_progress import BatchProgress, map_parallel

def warp(path: Path) -> Path:          # must be importable at module level
    ...

with BatchProgress(total=len(paths), description="Warping", workers=8) as progress:
    outputs = map_parallel(warp, paths, workers=8, progress=progress)

Results come back in input order. The pool is fed from a bounded in-flight window (2 * workers by default), so paths can be a generator over a directory tree with a million entries and neither the futures nor the file list are ever fully materialised.

Report what each worker is doing

from typer_rich_progress import BatchProgress, WorkerReporter, map_parallel

def warp(path: Path, reporter: WorkerReporter) -> Path:
    reporter.stage(f"reading {path.name}")
    data = read(path)
    reporter.stage(f"reprojecting {path.name}")
    return write(reproject(data))

with BatchProgress(total=len(paths), workers=8) as progress:
    map_parallel(warp, paths, workers=8, progress=progress, pass_reporter=True)
Warping  ━━━━━━━━━━╺━━━━━━━━━━━━ 217/620  35.0%  41.0s elapsed | ETA 1m16s | 5.29/s
Worker            Stage                            Done
ForkProcess-3     reprojecting n52e013.tif           28
ForkProcess-4     reading n52e014.tif                26
ForkProcess-5     writing n51e013_cog.tif            27

asyncio

from typer_rich_progress import BatchProgress, map_async

async def fetch(url: str) -> bytes:
    ...

with BatchProgress(total=len(urls), description="Fetching") as progress:
    blobs = await map_async(fetch, urls, concurrency=16, progress=progress)

Typer wiring

import typer
from typing import Annotated
from typer_rich_progress import map_parallel, typer_progress_options

app = typer.Typer()
opts = typer_progress_options()

@app.command()
def warp(
    src: Path,
    progress: Annotated[bool, opts.progress] = True,
    workers: Annotated[int, opts.workers] = 4,
) -> None:
    paths = sorted(src.glob("*.tif"))
    with opts.build(len(paths), "Warping", progress=progress, workers=workers) as bar:
        map_parallel(warp_one, paths, workers=workers, progress=bar)

That gives the command --progress/--no-progress and --workers/-j with validation, and nothing else has to change.

How it works

The ETA model

The estimator tracks an exponentially-weighted moving average of the wall-clock interval between completions arriving back at the parent, not the duration of individual worker calls. Using arrival intervals is what makes the number correct under parallelism: eight workers each taking 8 s produce one completion per second, and the interval already reflects that. Feeding it worker-side durations instead would over-estimate the remaining time by roughly the worker count.

Let d_n be the interval that item n took. The recurrence is the standard EWMA with zero initialisation and bias correction:

s_0    = 0
s_n    = α·d_n + (1 − α)·s_(n−1)
mean_n = s_n / (1 − (1 − α)^n)

ETA    = mean_n × remaining
rate   = 1 / mean_n

Two details matter.

The bias correction. Starting the sum at zero means s_1 = α·d_1, which for the default α = 0.15 is 15 % of the true first interval — the bar would open by claiming the job is nearly seven times faster than it is. Dividing by 1 − (1 − α)^n removes exactly that startup bias: at n = 1 the correction is α, so the estimate is d_1 precisely, and it decays to 1 as history accumulates. The alternative — seeding s_1 = d_1 — makes the first sample permanently over-weighted, which is worse for jobs whose first item is unrepresentative (a cold page cache, a JIT warm-up, an S3 connection being established).

The choice of α. The estimator's effective memory is about 1/α items, so the default α = 0.15 averages over roughly the last seven completions. Raise it towards 1 to track abrupt speed changes almost immediately at the cost of a jittery ETA; lower it towards 0 for a stable number that lags behind reality. What it buys you, concretely: after twenty 1-second items followed by five 10-second items, elapsed / completed still reports 2.8 s/item, while the EWMA at α = 0.5 reports 9.72 s/item. Only one of those is a useful thing to show a user waiting on the remaining 500 tiles.

Before the first completion there is genuinely nothing to extrapolate from, so the ETA is None and renders as -- rather than as an invented number.

Worker reporting and pickling

Only the parent process owns the terminal. Children never render — they put small frozen dataclasses on a queue, and the parent's renderer drains them without blocking on each refresh. That constrains what can cross the boundary, and the constraints are worth stating plainly:

  • The queue must be a multiprocessing.Manager().Queue() proxy. A plain multiprocessing.Queue cannot be pickled into a ProcessPoolExecutor task argument; manager proxies can. BatchProgress.event_queue() starts the manager lazily, so jobs that never report stages never pay for the extra process, and shuts it down on exit.
  • The function you pass to map_parallel must be importable at module level. Lambdas, closures and locally-defined functions are not picklable and will fail immediately.
  • Events carry only str and int, so no user object is dragged across the boundary and no unpicklable attribute can sneak in with it.
  • WorkerReporter.advance() is for sub-items inside one task. map_parallel already counts the task itself when its future resolves; calling advance() for the task as a whole double-counts it.

For asyncio there is no process boundary, so a plain queue.Queue is used instead and is attached automatically.

Bounded submission

Both maps keep at most a fixed number of items in flight and pull from the input iterable one item per submission. Beyond memory, this is what keeps the ETA meaningful: submitting everything at once would mark every item as started simultaneously, and the interval signal the estimator depends on would collapse.

A failure in any task cancels the outstanding window and raises WorkerTaskError, which names the index and the item and keeps the original exception as __cause__. Nothing is swallowed.

Degrading outside a terminal

Animation is disabled, in order, when: enabled=False is passed; NO_COLOR is set; CI is set; TERM is dumb or unknown; or the console is not a terminal. In that mode the display writes periodic plain-text lines straight to the console's underlying file, bypassing Rich entirely, so not one escape byte can reach a log:

Warping: starting, 620 items (CI is set)
Warping: 96/620 (15.5%) elapsed 18.2s eta 1m39s 5.27 items/s
Warping: 213/620 (34.4%) elapsed 40.1s eta 1m17s 5.31 items/s
Warping: 620/620 (100.0%) elapsed 1m56s eta 0.0s 5.33 items/s [done]

The counters, the summary and the exit behaviour are unchanged — --no-progress still tells you what happened, it just stops moving.

API reference

Object Purpose
BatchProgress(total, description, *, workers, console, enabled, env, time_source, alpha, plain_interval, show_workers) The display; a context manager.
BatchProgress.track(iterable) Drop-in rich.progress.track replacement that feeds the ETA model.
BatchProgress.advance(n, *, worker, stage) Record completions, optionally attributed to a worker.
BatchProgress.set_stage(worker, stage) Update a worker's label without touching counters.
BatchProgress.reporter(worker) / .event_queue() / .drain() Worker-event plumbing.
BatchProgress.summary() The current state as one escape-free line.
map_parallel(fn, items, *, workers, progress, max_in_flight, pass_reporter, executor) Bounded ProcessPoolExecutor map.
map_async(coro_fn, items, *, concurrency, progress, pass_reporter) Bounded asyncio map.
WorkerReporter(queue, worker) Worker-side stage() / advance() handle.
EtaEstimator(alpha, start) The EWMA model, usable on its own.
resolve_animation(console, *, enabled, env) The TTY/CI/NO_COLOR decision, with a reason.
typer_progress_options(...) Reusable --progress/--no-progress and --workers options.

Two arguments exist mainly for testing and are worth knowing about: time_source injects a clock, so ETA and throughput can be asserted exactly with no sleeps, and env injects the environment mapping consulted by the TTY detection — pass env={} so a surrounding CI variable cannot change the outcome under test.

Development

uv sync --all-extras
uv run ruff check .
uv run ruff format --check .
uv run mypy
uv run pytest

pytest is configured to fail under 85 % statement coverage on src/. The suite drives the clock manually rather than sleeping, renders into a StringIO console and asserts on the captured output, so it is fast and does not flake.

Further reading

The design decisions here are discussed at greater length in these guides:

License

MIT. See LICENSE.

About

Consistent Rich progress bars, smoothed ETA and per-worker status for any multiprocessing or asyncio batch job

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages