Warm-up · Activity 1 of 7
// A3.3 · ~30 min · Advanced
concurrent.futures
After this lesson you can hand work to a thread or process pool with one interface, collect results as they finish, and handle the exceptions that come back.
You will be able to
- Submit calls to ThreadPoolExecutor and ProcessPoolExecutor and collect results with result(), map and as_completed
- Handle exceptions from futures with result(), exception() and try and except
- Choose threads for I/O-bound and processes or interpreters for CPU-bound work
Predict · Activity 2 of 7
Predict before you read on. The submitted call fails. What does this print?
from concurrent.futures import ThreadPoolExecutor with ThreadPoolExecutor(max_workers=2) as executor: future = executor.submit(int, "forty-two") try: print(future.result()) except ValueError as error: print("result() raised:", error)Practice · Activity 3 of 7
Fill in the executor method that schedules pow(2, 10) and returns a Future.
from concurrent.futures import ThreadPoolExecutor with ThreadPoolExecutor() as executor: future = executor.____(pow, 2, 10) print(future.result())executor.(pow, 2, 10)Practice · Activity 4 of 7
exception() returns what the call raised. What does this print?
from concurrent.futures import ThreadPoolExecutor with ThreadPoolExecutor() as executor: ok = executor.submit(int, "7") bad = executor.submit(int, "seven") print(ok.exception(), repr(bad.exception()), ok.done(), bad.done())Practice · Activity 5 of 7
Match each call to what it gives you.
Brain teaser · Activity 6 of 7
Brain teaser. The submitted call raises ValueError, and nobody asks for its result. What does this print?
from concurrent.futures import ThreadPoolExecutor with ThreadPoolExecutor() as executor: executor.submit(int, "seven") print("finished")Apply · Activity 7 of 7
Mini-task, on your own computer. With a ThreadPoolExecutor of three workers, submit len(word) for each of "kiwi", "fig" and "banana". Keep a dict from each future to its word, collect the results with as_completed into a dict of word to length, and print its sorted items.
Check your work against this list
Build it yourself
Read the worked example, then write the exercises. Your code runs in your browser or on your computer and is never uploaded.
Worked example
Fetch in threads, count in processes
Four page fetches are I/O-bound, so they run in a ThreadPoolExecutor and are collected with as_completed; the missing page raises KeyError, which result() brings out. Counting vowels is CPU work, so it runs in a ProcessPoolExecutor with map, which keeps the input order. The worker functions live in textwork.py, and a pretend server with sleep keeps the example off the network. Save both files and run python main.py (python3 main.py on macOS and Linux) on your own computer.
main.py
from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor, as_completed
from textwork import PAGES, count_vowels, fetch
def main() -> None:
# I/O-bound: fetching pages overlaps well in threads.
names = ["home", "about", "missing", "blog"]
pages: dict[str, str] = {}
with ThreadPoolExecutor(max_workers=4) as executor:
futures = {executor.submit(fetch, name): name for name in names}
for future in as_completed(futures): # in the order they finish
name = futures[future]
try:
pages[name] = future.result()
except KeyError as error:
print("could not fetch:", error)
print("fetched:", sorted(pages))
# CPU-bound: counting runs in worker processes, results in input order.
with ProcessPoolExecutor(max_workers=2) as executor:
counts = executor.map(count_vowels, [pages[name] for name in sorted(pages)])
for name, count in zip(sorted(pages), counts):
print(f"{name}: {count} vowels")
print("pages on the server:", len(PAGES))
if __name__ == "__main__":
main()
textwork.py
import time
# A pretend web server, so the example needs no network.
PAGES = {
"home": "welcome to the example site",
"about": "we write about python",
"blog": "threads processes and asyncio",
}
def fetch(name: str) -> str:
"""Pretend to download a page: wait a little, like network I/O."""
time.sleep(0.01)
return PAGES[name] # KeyError for a page that does not exist
def count_vowels(text: str) -> int:
"""CPU work: count the vowels in text."""
return sum(1 for char in text if char in "aeiou")
Run it with
python main.pyOutput
could not fetch: 'missing'
fetched: ['about', 'blog', 'home']
about: 7 vowels
blog: 9 vowels
home: 10 vowels
pages on the server: 3- The KeyError happened in a pool thread but was raised by future.result() in the main thread.
- sorted(pages) makes the output the same whichever fetch finished first.
- executor.map returned the counts in the order of its input, so zip pairs each name with its count.
- Both executors sit in with blocks, so each waits for its work before the program goes on.
Exercises
Exercise 1 of 2
Fetch everything, keep the errors
Complete fetch_all(names, fetch, workers). Submit fetch(name) for every name to a ThreadPoolExecutor with workers threads, collect the futures with as_completed, and return a dict from each name to its text. A fetch that raises must not stop the others: store "error: " followed by the exception for that name. Run the tests on your own computer.
This exercise needs Python on your computer (the browser version cannot run it). The files and commands are below.
Hints
Hint 1
Build futures = {executor.submit(fetch, name): name for name in names} inside with ThreadPoolExecutor(max_workers=workers) as executor:.
Hint 2
for future in as_completed(futures): gives each future when it is done; futures[future] is its name.
Hint 3
Call future.result() inside try; in except Exception as error: store f"error: {error}". str() of KeyError('missing') is 'missing' with quotes.
Show a solution
One way to solve it. Yours can look different and still pass the checks.
from collections.abc import Callable
from concurrent.futures import ThreadPoolExecutor, as_completed
def fetch_all(names: list[str], fetch: Callable[[str], str], workers: int = 4) -> dict[str, str]:
"""Fetch every name in a thread pool. Map each name to its text, or to "error: <exception>"."""
results: dict[str, str] = {}
with ThreadPoolExecutor(max_workers=workers) as executor:
futures = {executor.submit(fetch, name): name for name in names}
for future in as_completed(futures):
name = futures[future]
try:
results[name] = future.result()
except Exception as error:
results[name] = f"error: {error}"
return results
Run it on your computer
Install Python 3.14 or newer. Save these files in one folder, open a terminal in that folder, and run the commands below.
main.py
from collections.abc import Callable
from concurrent.futures import ThreadPoolExecutor, as_completed
def fetch_all(names: list[str], fetch: Callable[[str], str], workers: int = 4) -> dict[str, str]:
"""Fetch every name in a thread pool. Map each name to its text, or to "error: <exception>"."""
# One after another, and the first error stops everything.
# Submit every fetch to a ThreadPoolExecutor and collect them with as_completed.
return {name: fetch(name) for name in names}
test_main.py
import threading
from main import fetch_all
PAGES = {"home": "welcome", "about": "about us"}
def fake_fetch(name):
return PAGES[name]
def test_all_pages():
"""Every page is fetched and stored under its name"""
got = fetch_all(["home", "about"], fake_fetch)
assert got == {"home": "welcome", "about": "about us"}, f"fetch_all returned {got!r}"
def test_errors_are_kept():
"""A failing fetch gives "error: <exception>" and the others still arrive"""
got = fetch_all(["home", "missing"], fake_fetch)
assert got == {"home": "welcome", "missing": "error: 'missing'"}, f"fetch_all returned {got!r}"
def test_fetches_overlap():
"""Three fetches run at the same time in pool threads"""
barrier = threading.Barrier(3, timeout=0.5)
def slow_fetch(name):
barrier.wait()
return threading.current_thread().name
try:
got = fetch_all(["a", "b", "c"], slow_fetch, workers=3)
except threading.BrokenBarrierError:
got = {}
threads = sorted(got.values())
assert len(threads) == 3 and "MainThread" not in threads, f"the fetches ran in {threads!r}: submit them all to the executor"
On macOS and Linux, type python3 wherever these commands say python, as in the first lesson.
Run the program:
python main.pyRun the checks (needs learnrun.py in the same folder):
python learnrun.py testDownload learnrun.pyExercise 2 of 2
Sort out successes and failures
run_checks submits check_positive for each number to a ProcessPoolExecutor, but result() raises the ValueError of a negative number and ends the whole function. Change the loop so that it asks each future for its exception() first: collect successful results in results and the messages of failures, str(error), in errors. Return both lists sorted. Run the tests on your own computer.
This exercise needs Python on your computer (the browser version cannot run it). The files and commands are below.
Hints
Hint 1
future.exception() waits for the call like result() does, but returns the exception instead of raising it.
Hint 2
If it returns None, the call succeeded and future.result() is safe; otherwise append str(error) to errors.
Hint 3
A try and except ValueError around future.result() works as well.
Show a solution
One way to solve it. Yours can look different and still pass the checks.
from concurrent.futures import ProcessPoolExecutor
def check_positive(n: int) -> int:
"""Return n squared, or raise ValueError for a negative n."""
if n < 0:
raise ValueError(f"{n} is negative")
return n * n
def run_checks(numbers: list[int]) -> tuple[list[int], list[str]]:
"""Run check_positive in worker processes; return (sorted results, sorted error messages)."""
results: list[int] = []
errors: list[str] = []
with ProcessPoolExecutor(max_workers=2) as executor:
futures = [executor.submit(check_positive, n) for n in numbers]
for future in futures:
error = future.exception() # waits; None when the call succeeded
if error is None:
results.append(future.result())
else:
errors.append(str(error))
return sorted(results), sorted(errors)
if __name__ == "__main__":
print(run_checks([3, -1, 2]))
Run it on your computer
Install Python 3.14 or newer. Save these files in one folder, open a terminal in that folder, and run the commands below.
main.py
from concurrent.futures import ProcessPoolExecutor
def check_positive(n: int) -> int:
"""Return n squared, or raise ValueError for a negative n."""
if n < 0:
raise ValueError(f"{n} is negative")
return n * n
def run_checks(numbers: list[int]) -> tuple[list[int], list[str]]:
"""Run check_positive in worker processes; return (sorted results, sorted error messages)."""
results: list[int] = []
errors: list[str] = []
with ProcessPoolExecutor(max_workers=2) as executor:
futures = [executor.submit(check_positive, n) for n in numbers]
for future in futures:
# result() raises the worker's ValueError here. Collect it in errors instead.
results.append(future.result())
return sorted(results), sorted(errors)
if __name__ == "__main__":
print(run_checks([3, -1, 2]))
test_main.py
from main import run_checks
def test_all_good():
"""Positive numbers give their sorted squares and no errors"""
got = run_checks([3, 1, 2])
assert got == ([1, 4, 9], []), f"run_checks([3, 1, 2]) returned {got!r}"
def test_errors_collected():
"""Negative numbers end up as error messages; the others still count"""
got = run_checks([3, -1, 2, -5])
assert got == ([4, 9], ["-1 is negative", "-5 is negative"]), f"run_checks([3, -1, 2, -5]) returned {got!r}"
def test_only_errors():
"""A list of negatives gives no results and one message each"""
got = run_checks([-2])
assert got == ([], ["-2 is negative"]), f"run_checks([-2]) returned {got!r}"
On macOS and Linux, type python3 wherever these commands say python, as in the first lesson.
Run the program:
python main.pyRun the checks (needs learnrun.py in the same folder):
python learnrun.py testDownload learnrun.pyCommon mistakes
Calling the function instead of submitting it
from concurrent.futures import ThreadPoolExecutor
def square(x: int) -> int:
return x * x
with ThreadPoolExecutor() as executor:
future = executor.submit(square(3))
print(future.result())
What Python prints
TypeError: 'int' object is not callableWhy, and the fix
square(3) runs right away in the main thread, and submit receives its result, 9. The worker then tries to call 9(). Pass the function and its arguments separately: executor.submit(square, 3). The error only shows up when you call result().
Submitting after the with block
from concurrent.futures import ThreadPoolExecutor
with ThreadPoolExecutor() as executor:
first = executor.submit(pow, 2, 10)
print(first.result())
second = executor.submit(pow, 3, 3)
What Python prints
RuntimeError: cannot schedule new futures after shutdownWhy, and the fix
Leaving the with block calls shutdown(wait=True): the executor finishes its work and accepts nothing new. Futures you already have still work, so first.result() prints 1024. Submit everything inside the with block.
Expecting map to skip a failed item
from concurrent.futures import ThreadPoolExecutor
with ThreadPoolExecutor() as executor:
for n in executor.map(int, ["1", "two", "3"]):
print(n)
What Python prints
ValueError: invalid literal for int() with base 10: 'two'Why, and the fix
executor.map raises a call's exception when the loop reaches that result, so 1 is printed and then the loop stops: 3 is never printed. To keep going past failures, submit each call and handle result() in try and except per future.
Python in the browser: Pyodide 314.0.7, MPL-2.0. Licence and source
Exit ticket
5 questions, no hints. Score 80% or more to complete the lesson.
Finish every activity above to unlock the exit ticket.