Zum Inhalt springen
aviral gupta

// A3.3 · ca. 30 Min. · Vertiefung

concurrent.futures

Nach dieser Lektion geben Sie Arbeit über eine einzige Schnittstelle an einen Thread- oder Prozess-Pool, sammeln Ergebnisse ein, sobald sie fertig sind, und behandeln die Exceptions, die zurückkommen.

Lektion 3 von 6 in A3 Nebenläufigkeit

Danach können Sie

  • Aufrufe an ThreadPoolExecutor und ProcessPoolExecutor übergeben und Ergebnisse mit result(), map und as_completed einsammeln
  • Exceptions aus Futures mit result(), exception() und try und except behandeln
  • Threads für I/O-lastige und Prozesse oder Interpreter für CPU-lastige Arbeit wählen
  1. Aufwärmen · Aufgabe 1 von 7

    Aufwärmen aus der letzten Lektion: Warum können Sie multiprocessing.Pool.map kein lambda übergeben?

  2. Vorhersagen · Aufgabe 2 von 7

    Sagen Sie es vorher, bevor Sie weiterlesen. Der eingereichte Aufruf scheitert. Was gibt das aus?

    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)
  3. Üben · Aufgabe 3 von 7

    Setzen Sie die Executor-Methode ein, die pow(2, 10) einplant und ein Future liefert.

    from concurrent.futures import ThreadPoolExecutor
    
    with ThreadPoolExecutor() as executor:
        future = executor.____(pow, 2, 10)
    print(future.result())
    executor.(pow, 2, 10)
  4. Üben · Aufgabe 4 von 7

    exception() liefert, was der Aufruf ausgelöst hat. Was gibt das aus?

    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())
  5. Üben · Aufgabe 5 von 7

    Ordnen Sie jedem Aufruf zu, was er Ihnen liefert.

  6. Denksport · Aufgabe 6 von 7

    Knobelaufgabe. Der eingereichte Aufruf löst ValueError aus, und niemand fragt nach seinem Ergebnis. Was gibt das aus?

    from concurrent.futures import ThreadPoolExecutor
    
    with ThreadPoolExecutor() as executor:
        executor.submit(int, "seven")
    print("finished")
  7. Anwenden · Aufgabe 7 von 7

    Mini-Aufgabe, auf Ihrem eigenen Rechner. Reichen Sie mit einem ThreadPoolExecutor mit drei Workern len(word) für "kiwi", "fig" und "banana" ein. Führen Sie ein Dict von jedem Future zu seinem Wort, sammeln Sie die Ergebnisse mit as_completed in einem Dict von Wort zu Länge und geben Sie dessen sortierte Einträge aus.

    Prüfen Sie Ihr Ergebnis anhand dieser Liste

Selbst programmieren

Lesen Sie das ausgearbeitete Beispiel und lösen Sie dann die Übungen. Ihr Code läuft in Ihrem Browser oder auf Ihrem Computer und wird nie hochgeladen.

Ausgearbeitetes Beispiel

In Threads laden, in Prozessen zählen

Vier Seitenabrufe sind I/O-lastig, also laufen sie in einem ThreadPoolExecutor und werden mit as_completed eingesammelt; die fehlende Seite löst KeyError aus, den result() hervorholt. Vokale zählen ist CPU-Arbeit, also läuft es in einem ProcessPoolExecutor mit map, das die Reihenfolge der Eingabe behält. Die Worker-Funktionen stehen in textwork.py, und ein gespielter Server mit sleep hält das Beispiel vom Netzwerk fern. Speichern Sie beide Dateien und starten Sie python main.py (unter macOS und Linux python3 main.py) auf Ihrem Rechner.

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")

Ausführen mit

python main.py

Ausgabe

could not fetch: 'missing'
fetched: ['about', 'blog', 'home']
about: 7 vowels
blog: 9 vowels
home: 10 vowels
pages on the server: 3
  • Der KeyError entstand in einem Pool-Thread, wurde aber von future.result() im Haupt-Thread ausgelöst.
  • sorted(pages) macht die Ausgabe gleich, egal welcher Abruf zuerst fertig war.
  • executor.map lieferte die Anzahlen in der Reihenfolge seiner Eingabe, also passt zip jeden Namen zu seiner Anzahl.
  • Beide Executors stehen in with-Blöcken, also wartet jeder auf seine Arbeit, bevor das Programm weitergeht.

Übungen

Übung 1 von 2

Alles laden, Fehler behalten

Vervollständigen Sie fetch_all(names, fetch, workers). Reichen Sie fetch(name) für jeden Namen bei einem ThreadPoolExecutor mit workers Threads ein, sammeln Sie die Futures mit as_completed ein und liefern Sie ein Dict von jedem Namen zu seinem Text. Ein Abruf, der auslöst, darf die anderen nicht stoppen: Speichern Sie für diesen Namen "error: " gefolgt von der Exception. Führen Sie die Tests auf Ihrem Rechner aus.

Diese Übung braucht Python auf Ihrem Computer (die Browser-Version kann sie nicht ausführen). Dateien und Befehle stehen unten.

Hinweise
  1. Hinweis 1

    Bauen Sie futures = {executor.submit(fetch, name): name for name in names} innerhalb von with ThreadPoolExecutor(max_workers=workers) as executor:.

  2. Hinweis 2

    for future in as_completed(futures): liefert jedes Future, sobald es fertig ist; futures[future] ist sein Name.

  3. Hinweis 3

    Rufen Sie future.result() in try auf; in except Exception as error: speichern Sie f"error: {error}". str() von KeyError('missing') ist 'missing' mit Anführungszeichen.

Eine Lösung zeigen

Ein möglicher Lösungsweg. Ihrer kann anders aussehen und trotzdem alle Prüfungen bestehen.

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
Auf dem eigenen Computer ausführen

Installieren Sie Python 3.14 oder neuer. Speichern Sie diese Dateien in einem Ordner, öffnen Sie dort ein Terminal und führen Sie die Befehle unten aus.

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():
    """Jede Seite wird geladen und unter ihrem Namen gespeichert"""
    got = fetch_all(["home", "about"], fake_fetch)
    assert got == {"home": "welcome", "about": "about us"}, f"fetch_all lieferte {got!r}"


def test_errors_are_kept():
    """Ein gescheiterter Abruf ergibt "error: <Exception>", die anderen kommen trotzdem an"""
    got = fetch_all(["home", "missing"], fake_fetch)
    assert got == {"home": "welcome", "missing": "error: 'missing'"}, f"fetch_all lieferte {got!r}"


def test_fetches_overlap():
    """Drei Abrufe laufen gleichzeitig 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"die Abrufe liefen in {threads!r}: reichen Sie alle beim Executor ein"

Unter macOS und Linux tippen Sie python3, wo in diesen Befehlen python steht, wie in der ersten Lektion.

Programm ausführen:

python main.py

Prüfungen ausführen (learnrun.py muss im selben Ordner liegen):

python learnrun.py test
learnrun.py herunterladen

Übung 2 von 2

Erfolge und Fehlschläge trennen

run_checks reicht check_positive für jede Zahl bei einem ProcessPoolExecutor ein, aber result() löst bei einer negativen Zahl den ValueError aus und beendet die ganze Funktion. Ändern Sie die Schleife so, dass sie jedes Future zuerst nach exception() fragt: Erfolgreiche Ergebnisse kommen in results, die Meldungen der Fehlschläge, str(error), in errors. Liefern Sie beide Listen sortiert. Führen Sie die Tests auf Ihrem Rechner aus.

Diese Übung braucht Python auf Ihrem Computer (die Browser-Version kann sie nicht ausführen). Dateien und Befehle stehen unten.

Hinweise
  1. Hinweis 1

    future.exception() wartet wie result() auf den Aufruf, liefert die Exception aber, statt sie auszulösen.

  2. Hinweis 2

    Liefert es None, war der Aufruf erfolgreich und future.result() ist sicher; sonst hängen Sie str(error) an errors an.

  3. Hinweis 3

    Ein try und except ValueError um future.result() funktioniert ebenfalls.

Eine Lösung zeigen

Ein möglicher Lösungsweg. Ihrer kann anders aussehen und trotzdem alle Prüfungen bestehen.

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]))
Auf dem eigenen Computer ausführen

Installieren Sie Python 3.14 oder neuer. Speichern Sie diese Dateien in einem Ordner, öffnen Sie dort ein Terminal und führen Sie die Befehle unten aus.

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 Zahlen ergeben ihre sortierten Quadrate und keine Fehler"""
    got = run_checks([3, 1, 2])
    assert got == ([1, 4, 9], []), f"run_checks([3, 1, 2]) lieferte {got!r}"


def test_errors_collected():
    """Negative Zahlen landen als Fehlermeldungen; die anderen zählen trotzdem"""
    got = run_checks([3, -1, 2, -5])
    assert got == ([4, 9], ["-1 is negative", "-5 is negative"]), f"run_checks([3, -1, 2, -5]) lieferte {got!r}"


def test_only_errors():
    """Nur negative Zahlen ergeben keine Ergebnisse und je eine Meldung"""
    got = run_checks([-2])
    assert got == ([], ["-2 is negative"]), f"run_checks([-2]) lieferte {got!r}"

Unter macOS und Linux tippen Sie python3, wo in diesen Befehlen python steht, wie in der ersten Lektion.

Programm ausführen:

python main.py

Prüfungen ausführen (learnrun.py muss im selben Ordner liegen):

python learnrun.py test
learnrun.py herunterladen

Häufige Fehler

Die Funktion aufrufen statt sie einzureichen

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())

Was Python ausgibt

TypeError: 'int' object is not callable

Warum, und die Lösung

square(3) läuft sofort im Haupt-Thread, und submit bekommt sein Ergebnis, 9. Der Worker versucht dann, 9() aufzurufen. Übergeben Sie Funktion und Argumente getrennt: executor.submit(square, 3). Der Fehler zeigt sich erst, wenn Sie result() aufrufen.

Nach dem with-Block einreichen

from concurrent.futures import ThreadPoolExecutor

with ThreadPoolExecutor() as executor:
    first = executor.submit(pow, 2, 10)
print(first.result())
second = executor.submit(pow, 3, 3)

Was Python ausgibt

RuntimeError: cannot schedule new futures after shutdown

Warum, und die Lösung

Das Verlassen des with-Blocks ruft shutdown(wait=True) auf: Der Executor beendet seine Arbeit und nimmt nichts Neues mehr an. Futures, die Sie schon haben, funktionieren weiter, first.result() gibt also 1024 aus. Reichen Sie alles innerhalb des with-Blocks ein.

Erwarten, dass map ein gescheitertes Element überspringt

from concurrent.futures import ThreadPoolExecutor

with ThreadPoolExecutor() as executor:
    for n in executor.map(int, ["1", "two", "3"]):
        print(n)

Was Python ausgibt

ValueError: invalid literal for int() with base 10: 'two'

Warum, und die Lösung

executor.map löst die Exception eines Aufrufs aus, wenn die Schleife dieses Ergebnis erreicht: 1 wird ausgegeben, dann endet die Schleife, und 3 erscheint nie. Um über Fehlschläge hinweg weiterzumachen, reichen Sie jeden Aufruf einzeln ein und behandeln result() pro Future in try und except.

Python im Browser: Pyodide 314.0.7, MPL-2.0. Lizenz und Quellcode

Abschlussquiz

5 Fragen, ohne Hinweise. Ab 80 % ist die Lektion abgeschlossen.

Erledigen Sie zuerst alle Aufgaben oben, um das Abschlussquiz freizuschalten.

Problem melden

Etwas ist falsch oder unklar? Beschreiben Sie es kurz, dann wird es geprüft und korrigiert.

#

Mindestens 20 Zeichen.

Nur, wenn Sie eine Antwort wünschen.

Kernideen

Executors und Futures

concurrent.futures bietet Threads und Prozessen eine gemeinsame, abstrakte Schnittstelle. executor.submit(fn, *args) plant fn(*args) ein und liefert sofort ein Future. future.result() wartet auf den Aufruf und liefert seinen Wert. executor.map(fn, items) arbeitet wie map und liefert die Ergebnisse in der Reihenfolge von items. Nutzen Sie den Executor in einem with-Block: Am Ende fährt er herunter und wartet auf alle offenen Aufrufe. as_completed(futures) liefert Futures in der Reihenfolge, in der sie fertig werden; führen Sie deshalb ein Dict vom Future zur Eingabe, etwa {executor.submit(fetch, n): n for n in names}.

Exceptions reisen mit dem Future

Eine Exception in einem eingereichten Aufruf taucht nicht dort auf, wo Sie submit() aufrufen. Sie wird im Future gespeichert: result() löst dieselbe Exception aus, und exception() liefert sie, nach einem Erfolg None. Bei map wird die Exception ausgelöst, wenn Sie dieses Element in den Ergebnissen erreichen, und die weitere Iteration endet. Fragen Sie ein Future nie nach seinem Ergebnis, sieht niemand seine Exception: Rufen Sie auf jedem Future result() oder exception() auf, in try und except, wo etwas scheitern kann.

Welcher Executor?

ThreadPoolExecutor passt zu I/O-lastiger Arbeit; sein Standard für max_workers ist min(32, (os.process_cpu_count() or 1) + 4). ProcessPoolExecutor passt zu CPU-lastiger Arbeit und bringt die Regeln der letzten Lektion mit: nur picklebare Funktionen, Argumente und Ergebnisse, Funktionen auf Modulebene und ein Hauptmodul, das die Worker importieren können, also den __main__-Guard. In 3.14 gibt es außerdem InterpreterPoolExecutor: Jeder Worker-Thread führt einen eigenen Interpreter mit eigenem GIL aus, rechnet also parallel, und pickelt Aufrufe und Ergebnisse ebenfalls.

Quellen

Zuletzt geprüft am 29. September 2026