skip to content

How do you stop queued concurrent.futures work when a 6-hour nightly billing run hits a fatal error?

level: seniorimportance: should knowfreq 40%

answer

  1. only one state can still be cancelled
  2. the with block does not abort anything
  3. one keyword drops the whole backlog
  4. running items must agree to stop
  5. stopping work is not undoing work

basics

~20 s

Call Executor.shutdown(wait=True, cancel_futures=True): it drops everything still queued and waits for the jobs already running. Future.cancel() only works while a job is still pending, so anything in flight has to finish or be stopped cooperatively.

solid answer

~50 s

`Future.cancel()` succeeds only while the job is still **pending**; once a worker has picked it up it returns `False`, and `concurrent.futures` has no way to interrupt a running callable. So stopping a long run has two halves. For queued work, call `Executor.shutdown(wait=True, cancel_futures=True)` — it cancels every not-yet-started `Future` and then waits for the in-flight ones. Note that exiting the `with` block does **not** do this: `Executor.__exit__` calls `shutdown(wait=True)` with `cancel_futures=False`, which drains the entire backlog. For running work you need cooperation: a `threading.Event` the workers poll between steps, so they return early instead of being killed. To notice the failure at all, submit yourself and use `concurrent.futures.wait(futures, return_when=FIRST_EXCEPTION)` rather than draining results in order. Cancellation is not a rollback — items already committed stay committed, so the run needs a per-item ledger and idempotent retries.

code

python · 15 lines
python
import time
from concurrent.futures import ThreadPoolExecutor, wait, FIRST_EXCEPTION

def charge(n):
    time.sleep(0.05)
    if n == 2:
        raise RuntimeError("gateway rejected the batch")
    return n

pool = ThreadPoolExecutor(max_workers=2)
futures = [pool.submit(charge, n) for n in range(40)]
done, not_done = wait(futures, return_when=FIRST_EXCEPTION)
pool.shutdown(wait=True, cancel_futures=True)
print("still queued when we stopped:", len(not_done))
print("cancelled:", sum(f.cancelled() for f in futures))

go deeper

for a junior

Know the two states that matter: a pending job can be cancelled, a running one cannot. Remember that shutdown takes wait and cancel_futures, and that leaving the with block waits for everything.

for a middle

Explain the mechanics: Future.cancel() returns False once a worker has the item, shutdown(cancel_futures=True) drains the pending queue, and __exit__ does not pass that flag, so an early exit still processes the backlog.

for a senior

Demonstrate the whole abort path on a long batch: detect the failure early with FIRST_EXCEPTION or a failure-rate threshold, drop the backlog, signal running items through an event, and report succeeded/failed/cancelled counts.

for a principal

Own the recovery contract, not the API: per-item idempotency keys and a durable outcome ledger so a half-finished run is resumable, an explicit unit of atomicity, and an abort policy the on-call engineer can predict at three in the morning.

A nightly subscription-billing run that fans 200,000 accounts across a pool and takes six hours has to answer a hard question halfway through: the payment gateway starts rejecting everything, and you want the run to stop **now** rather than push through three more hours of guaranteed failures. `concurrent.futures` gives you exactly two levers, and it is important to know which one does what. ### Lever one: cancelling what has not started A `Future` moves through pending, running, then finished or cancelled. `Future.cancel()` returns `True` only from the pending state — it flips the future to cancelled and the worker skips it. Once a worker has picked the item up, `cancel()` returns `False` and does nothing. There is no `Future.kill()`, and there cannot be: interrupting arbitrary Python code at an arbitrary bytecode boundary would leave half-updated state behind. `Executor.shutdown(wait=True, cancel_futures=True)` is the bulk form of the same thing. It drains the pending work queue, cancels each of those futures, refuses new `submit` calls with `RuntimeError`, and then — because `wait=True` — blocks until the currently running items finish. `cancel_futures` was added in Python 3.9; before that you had to loop over your futures calling `cancel()` yourself, which is still a fine way to cancel a subset. The trap is the context manager. `Executor.__exit__` calls `shutdown(wait=True)` and nothing else, so leaving the `with` block on an exception **still drains the entire backlog** — the run you thought you aborted keeps billing for three more hours before the traceback prints. If a `with` block might be exited early, either construct the executor outside it and shut it down explicitly, or wrap the body in `try/except` and call `shutdown(cancel_futures=True)` on the way out. ### Lever two: stopping what is already running Running items can only be stopped cooperatively. The standard shape is a `threading.Event` created in the parent, closed over (thread pool) or passed via an `initializer` (process pool), and checked by the worker between units of work: ```python abort = threading.Event() def charge(account): if abort.is_set(): return "skipped" ... ``` That gives a bounded stop time equal to one item, not one queue. With a `ProcessPoolExecutor` the only forcible alternative is killing the child processes, which breaks the pool: every pending `Future` completes with a `BrokenProcessPool` error and the executor is finished. ### Noticing the failure early enough to act None of this helps if you learn about the failure last. Draining results with `Executor.map` or by iterating your futures in submission order means the first error you see is the first *by input position*, possibly hours late. Two better shapes: * `done, not_done = concurrent.futures.wait(futures, return_when=FIRST_EXCEPTION)` returns the moment any job raises, without raising itself, and hands you both sets. Follow it with `shutdown(wait=True, cancel_futures=True)`. * `for fut in concurrent.futures.as_completed(futures)` with a failure counter, so you can apply a policy — abort after N consecutive failures, or after a failure rate crosses a threshold — instead of aborting on the first blip. On a six-hour billing run a threshold is usually right: one declined card is normal, a 100% decline rate is an outage. ### Cancellation is not a rollback This is the part that separates a senior answer. Cancelling futures stops *future* work; it undoes nothing. The accounts already charged stay charged, and the executor has no notion of a transaction spanning items. A partial-failure rollback has to be designed at the job level: * Make each item idempotent and keyed — an idempotency key per invoice, so a re-run charges nobody twice. * Record each item's outcome durably as it completes, not at the end, so an aborted run leaves a ledger you can resume from rather than an all-or-nothing mystery. * Decide explicitly whether the unit of atomicity is the item or the batch. Per-item is almost always the answer for a long run, because a six-hour batch that must roll back wholesale has a six-hour recovery. * Have the failure path report counts — succeeded, failed, cancelled, never-started — because "the run stopped" is not an operable statement. ### One shutdown detail worth knowing Since Python 3.9, `ThreadPoolExecutor` workers are ordinary non-daemon threads that interpreter shutdown joins. A process holding an executor with a long queue will therefore not exit until that queue drains, even on `SIGINT`. Explicitly calling `shutdown(wait=True, cancel_futures=True)` in your signal handling path is what makes Ctrl-C actually stop a batch job in reasonable time.

  • What does Future.result() do on a future that was cancelled before it ran?
    It raises `CancelledError` rather than returning a value or the job's own exception, and `Future.exception()` raises it too instead of returning it. So when you drain futures after an abort you have to treat cancellation as a distinct outcome: `cancelled()` is true, `done()` is true, and there is no result to read. Counting those separately is what gives you an accurate never-started number in the run report.
  • Why does shutdown(wait=False) not make the abort faster in practice?
    `wait=False` only means your thread returns immediately; it does not stop anything. The already-running items keep running, and since 3.9 the worker threads are non-daemon, so interpreter exit joins them anyway. Combined with `cancel_futures=True` it is useful when you want to keep doing other work while the pool winds down, but the wall-clock stop time is still bounded by the longest running item.
  • How would you apply a failure threshold rather than aborting on the very first error?
    Consume with `as_completed` and keep counters. Increment a failure count when `fut.exception()` is not `None`, and once the count or the failure rate crosses a policy threshold, call `shutdown(wait=True, cancel_futures=True)` and set the abort event for running items. That distinguishes normal per-item failures, which a billing run always has, from a systemic outage, which is the only thing worth stopping for.

saying these in an interview costs you the question

  • Believes Future.cancel() can stop an already running job
  • Thinks exiting the with block cancels queued work
  • Expects shutdown(wait=False) to make work stop sooner
  • Treats cancelling futures as rolling back completed work
  • Waits for results in submission order and notices failures late
  • Plans to kill worker processes and keep using the same pool

context