Nebenläufigkeit: Event Loop, Backpressure und Cancellation durch eine Kette

Track Konzepte · Nebenläufigkeit und Last · ca. 55 Min.

Was du aus Teil a brauchst

Eine Coroutine hält nur an einem await an, zwischen zwei await kommt kein anderer dran, und ein Thread, der wartet, blockiert nichts, wenn er den GIL freigibt. Eine gemeinsame Ressource begrenzt den Durchsatz, egal wie viele Worker es gibt. Falls das neu ist: Teil 1.

Worum es geht

Ein einziger blockierender Aufruf in einer async def hält den ganzen Server an. Ein Erzeuger, der schneller liefert als der Verbraucher schafft, füllt den Speicher oder verliert still Daten. Ein Timeout, der die Arbeit darunter nicht beendet, leakt Verbindungen. Das sind die drei Fehler, die in async-Code am häufigsten Ausfälle verursachen.

Am Ende dieses Teils kannst du:

  • vorhersagen, in welcher Reihenfolge Coroutinen laufen und was ein blockierender Aufruf anrichtet (Event Loop),
  • eine Bounded Queue mit Backpressure bauen, damit ein schneller Erzeuger den Speicher nicht füllt,
  • Cancellation und Timeouts durch eine Kette asynchroner Aufrufe führen, ohne dass etwas weiterläuft (strukturierte Nebenläufigkeit, structured concurrency).

asyncio läuft im Browser, und wir geben ihm eine virtuelle Uhr, damit sleep(5) keine fünf Sekunden kostet und jede Ausgabe bei jedem Lauf gleich ist.

Plane ehrlich 55 Minuten ein: etwa 20 Minuten Lesen, 35 Minuten Übungen (drei Stück). Eine gute Pause ist nach Übung 2.

Von JS/TS her gedacht

Konzept JS/TS Python
Abbruch AbortController, AbortSignal (du reichst das Signal selbst durch) task.cancel(), CancelledError (kommt automatisch an jedem await an)
Timeout AbortSignal.timeout(ms) asyncio.timeout(s), asyncio.wait_for
Mehrere Aufgaben starten und einsammeln Promise.all asyncio.gather, besser asyncio.TaskGroup
Gegendruck bei Streams stream.write() liefert false, dann auf drain warten (allgemeines Fachwissen) asyncio.Queue(maxsize=n): await q.put() wartet, wenn voll

Wie Abbruch durch eine Kette läuft, siehst du in TypeScript so. Die Kette heißt handler ruft service ruft db. Ausgeführt mit Node (die Datei als JavaScript ohne Typannotationen, Timeout 200 ms, die Datenbank bräuchte 2000 ms):

function schlafen(ms: number, signal: AbortSignal): Promise<void> {
  return new Promise((resolve, reject) => {
    const timer = setTimeout(resolve, ms);
    signal.addEventListener("abort", () => {
      clearTimeout(timer);
      reject(signal.reason);
    }, { once: true });
  });
}
async function db(ms: number, signal: AbortSignal) {
  try { await schlafen(ms, signal); }
  catch (e: any) { console.log("db: abgebrochen,", e.name); throw e; }
}
async function service(ms: number, signal: AbortSignal) {
  try { await db(ms, signal); }          // ohne "signal" hier bricht die Kette NICHT ab
  finally { console.log("service: finally"); }
}
async function handler(ms: number) {
  try {
    await service(ms, AbortSignal.timeout(200));
    console.log("handler: ok");
  } catch (e: any) { console.log("handler:", e.name); }
}
await handler(2000);
db: abgebrochen, TimeoutError
service: finally
handler: TimeoutError

Der Unterschied: In JS musst du das signal in jede Funktion der Kette hineinreichen, sonst läuft die Arbeit darunter weiter. In Python kommt die Abbruch-Ausnahme von selbst an jedem await an. Darum ist der Python-Weg einfacher, aber es gibt andere Fallen (Übung 3).

Konzept

Schritt 1: Der Event Loop, Blockierung und Thread-Pools

In asyncio läuft alles in einem Thread. Ein Aufruf, der blockiert (time.sleep, eine synchrone HTTP-Bibliothek, eine lange Rechenschleife), hält den ganzen Loop an: Kein Timer feuert, kein anderer Request kommt dran. Für die Beispiele hier brauchst du einen Event Loop mit virtueller Uhr. Er verhält sich wie der echte asyncio-Loop, nur dass sleep(5) die Uhr einfach weiterstellt. Die Details musst du nicht verstehen. Die Idee: Wenn alle Tasks warten, springt die Uhr direkt zum nächsten Timer, statt wirklich zu schlafen. Wichtig ist laufe(coroutine): Das führt eine Coroutine aus (wie asyncio.run, nur mit der virtuellen Uhr) und gibt ihr Ergebnis zurück. asyncio.get_running_loop().time() liefert dann die virtuelle Zeit in Sekunden. Wo eine Übung laufe braucht, ist es schon geladen.

Der Ticker zählt alle 0,01 Sekunden hoch. Wir lassen die virtuelle Uhr 0,1 Sekunden verstreichen, einmal mit asyncio.sleep (gibt die Kontrolle ab) und einmal mit dem blockierenden time.sleep:

Ausgabe:

time.sleep   : 0 Ticks
asyncio.sleep: 10 Ticks

Während time.sleep den Thread festhält, kommt kein einziger Ticker-Schlag durch. Ist die blockierende Funktion nicht änderbar (Bibliothek ohne async), lagerst du sie in einen Thread-Pool aus. Das geht im Browser nicht (keine Threads), darum hier lokal ausgeführt mit Python 3.13.9 und der echten Uhr:

async def ausgelagert():
    z = {"n": 0}
    t = asyncio.create_task(ticker(z))
    await asyncio.sleep(0)
    ergebnis = await asyncio.to_thread(langsame_bibliothek)   # läuft in einem Worker-Thread
    t.cancel()
    return ergebnis, z["n"]
blockierend: Ergebnis fertig | Ticker-Schläge in 0,2 s: 0
to_thread:   Ergebnis fertig | Ticker-Schläge in 0,2 s mehr als 10: True (genau: 17 )

(langsame_bibliothek ruft time.sleep(0.2) auf. Die genaue Zahl 17 hängt vom Rechner ab, nicht die Aussage “0 gegen viele”.) CPU-lastige Arbeit im Loop lagerst du dagegen in einen Prozess aus, dazu mehr in Teil 1, Schritt 4.

Schritt 2: Backpressure und Bounded Queues

Ein Erzeuger liefert 3 Elemente pro Tick, der Verbraucher schafft 2. Was passiert mit der Warteschlange dazwischen? Hier der Simulator, den du in Übung 2 benutzt. Queue.put meldet mit False, wenn die Queue voll ist und das Element nicht aufgenommen wurde. lauf lässt pro Tick erst den Erzeuger, dann den Verbraucher arbeiten.

Ausgabe:

limit=None: erzeugt=300 verbraucht=200 in Queue=100 max Queue=102 verloren=0
limit=10: erzeugt=300 verbraucht=200 in Queue=8 max Queue=10 verloren=92

Zwei schlechte Wege. Unbegrenzte Queue (limit=None): nichts geht verloren, aber sie wächst um 1 pro Tick auf 100 Elemente. Nach Little’s Law wartet das letzte Element 100 / 2 = 50 Ticks, und im echten System geht irgendwann der Speicher aus. Begrenzte Queue, Erzeuger ignoriert das Signal (limit=10): der Speicher ist sicher, aber 92 Elemente gehen still verloren, weil das Ergebnis von put ignoriert wird.

Backpressure (Gegendruck) ist der dritte Weg: Die volle Queue bremst den Erzeuger. Er hört auf zu produzieren, bis wieder Platz da ist. Die Last wandert so zur Quelle zurück, wo man entscheiden kann: warten, langsamer senden oder bewusst ablehnen (Load Shedding, Lektion 07). In asyncio ist das eingebaut: await q.put(x) wartet, wenn die Queue ihr maxsize erreicht hat. Ein Verbraucher, der pro Sekunde ein Element nimmt, ein Erzeuger mit 5 Elementen, maxsize=2 (virtuelle Uhr in Sekunden):

Ausgabe:

 0.0 erzeuger will 0 abgeben
 0.0 erzeuger hat 0 abgegeben
 0.0 erzeuger will 1 abgeben
 0.0 erzeuger hat 1 abgegeben
 0.0 erzeuger will 2 abgeben
 1.0   verbraucher nimmt 0
 1.0 erzeuger hat 2 abgegeben
 1.0 erzeuger will 3 abgeben
 2.0   verbraucher nimmt 1
 2.0 erzeuger hat 3 abgegeben
 2.0 erzeuger will 4 abgeben
 3.0   verbraucher nimmt 2
 3.0 erzeuger hat 4 abgegeben

Der Erzeuger gibt zwei Elemente sofort ab, bleibt bei Element 2 hängen und macht erst weiter, wenn der Verbraucher bei t = 1,0 Platz schafft. Der Erzeuger läuft jetzt im Takt des Verbrauchers. Gute Queue-Größen sind klein: Eine große Queue versteckt ein Problem nur länger und verlängert die Wartezeit jedes Elements.

Schritt 3: Cancellation und Timeouts durch eine Kette

Ein Timeout ist wertlos, wenn die Arbeit darunter weiterläuft. Sie hält dann Verbindungen, Speicher und Locks fest (ein Leak) und rechnet für einen Aufrufer, der längst weg ist. Strukturierte Nebenläufigkeit (structured concurrency) fordert deshalb: Jede Aufgabe hat einen Besitzer und endet, bevor der Besitzer endet. Läuft ein Timeout ab oder wird abgebrochen, wird die ganze Kette darunter beendet und aufgeräumt, und erst dann geht es oben weiter.

In asyncio heißt Abbruch: An der Stelle, an der die Coroutine gerade wartet, wird eine CancelledError ausgelöst. Sie läuft durch jede Schicht nach oben, jedes finally räumt auf. Dieselbe Kette wie oben in TypeScript:

Ausgabe:

  service: start
  db: start
  db: abgebrochen, räume auf
  service: finally
  handler: Timeout nach 2.0 s (virtuell)

Die Reihenfolge ist wichtig: Erst die innerste Schicht räumt auf, dann die nächste, erst zuletzt meldet der Handler den Timeout. asyncio.timeout ist eine Deadline für alles, was im async with passiert. Drei Regeln:

  1. CancelledError nie verschlucken. except BaseException oder except CancelledError ohne raise hält den Abbruch an. Aufräumen im finally.
  2. task.cancel() ist nur eine Bitte. Der Task bekommt die Ausnahme erst beim nächsten Schedulerlauf. Wer cancel() ruft und sofort weitermacht, hat noch nicht aufgeräumt. Entweder await task abwarten oder (besser) Tasks gar nicht von Hand starten, sondern in einer TaskGroup.
  3. Fehler weitergeben. Scheitert ein Task in einer TaskGroup, werden die Geschwister abgebrochen, und die Fehler kommen gesammelt als ExceptionGroup beim Besitzer an (except*):

Ausgabe:

  a fertig
  c abgebrochen
  gefangen: ['b kaputt'] bei t = 2.0

a war vor dem Fehler fertig, c wurde bei t = 2 abgebrochen, der Fehler von b kommt beim Besitzer an. Nichts läuft weiter. Mit rohem asyncio.create_task und asyncio.gather musst du das selbst bauen (und vergisst es). Übung 3 zeigt genau diesen Fehler.

Falle

  1. Blockierender Aufruf im Event Loop. requests.get, time.sleep, eine große Schleife in einer async def: alles andere steht. Auslagern (to_thread, Prozess) oder eine async-Bibliothek nehmen.
  2. Unbegrenzte Queue. Sie sieht gesund aus, bis der Speicher voll ist. Jede Queue braucht eine Obergrenze und eine Antwort auf die Frage, was bei “voll” passiert (warten, drosseln, ablehnen).
  3. Backpressure ignoriert. Das False von put oder die drain-Meldung zu überhören verwandelt Gegendruck in stillen Datenverlust.
  4. Timeout, der nichts aufräumt. Ein wait mit Timeout um einen Task lässt den Task weiterlaufen. Beim nächsten Aufruf kommt der nächste Leak dazu.
  5. CancelledError verschlucken. Danach reagiert der Task nicht mehr auf Abbruch, der Besitzer hängt.

Übungen

Übung 1: Reihenfolge im Event Loop vorhersagen (ca. 8 Min.)

Drei Coroutinen laufen mit asyncio.gather(mahlen(), bruehen(), milch()) in einem Event Loop (wie in Schritt 1). Jede schreibt in dieselbe Liste log. time.sleep blockiert, asyncio.sleep gibt die Kontrolle ab, auch asyncio.sleep(0) gibt sie einmal ab.

async def mahlen():
    log.append("mahlen:start")
    await asyncio.sleep(1)
    log.append("mahlen:fertig")

async def bruehen():
    log.append("bruehen:start")
    time.sleep(0.01)                 # blockierend
    log.append("bruehen:vorbereitet")
    await asyncio.sleep(0.5)
    log.append("bruehen:fertig")

async def milch():
    log.append("milch:start")
    await asyncio.sleep(0)
    log.append("milch:aufgeschaeumt")

Trage die sieben Einträge in der Reihenfolge ein, in der sie in log landen. Verwende genau diese Texte: "mahlen:start", "mahlen:fertig", "bruehen:start", "bruehen:vorbereitet", "bruehen:fertig", "milch:start", "milch:aufgeschaeumt". Der Event Loop startet die Coroutinen in der Reihenfolge der Argumente von gather.

Gehe die Coroutinen in der Startreihenfolge durch und notiere, wo jede wirklich anhält. Was passiert bei time.sleep, was bei sleep(0), und nach welcher Regel kommt als Nächstes jemand dran, wenn alle warten?

antwort = ["mahlen:start", "bruehen:start", "bruehen:vorbereitet", "milch:start",
           "milch:aufgeschaeumt", "bruehen:fertig", "mahlen:fertig"]
antwort

mahlen läuft bis zum await asyncio.sleep(1) und hält an. bruehen läuft los: time.sleep gibt die Kontrolle nicht ab, also steht bruehen:vorbereitet gleich nach bruehen:start, bevor irgendeine andere Coroutine drankommt, und erst am await asyncio.sleep(0.5) hält bruehen an. milch schreibt milch:start und gibt mit sleep(0) nur kurz ab. Der Loop macht sofort mit ihr weiter, weil kein Timer wartet, der früher dran wäre, und milch:aufgeschaeumt folgt. Danach läuft die virtuelle Uhr: Der kürzere Timer (0,5 s) feuert vor dem längeren (1 s), also bruehen:fertig vor mahlen:fertig.

Übung 2: Bounded Queue und Backpressure (ca. 12 Min.)

Der Simulator aus Schritt 2 (Queue, Quelle, lauf, Naiv) ist geladen. Schreibe die Klasse Erzeuger so, dass sie Backpressure beachtet:

  • Erzeuger(quelle): quelle() liefert das nächste Element und zählt es als erzeugt. Neue Elemente gibt es nur über quelle().
  • schritt(q) wird jeden Tick einmal aufgerufen. Der Erzeuger will bis zu 3 Elemente abgeben (q.put(element)).
  • Ist die Queue voll (put liefert False), wurde das Element nicht aufgenommen. Dann darfst du es nicht verlieren, du erzeugst nichts Neues und machst im nächsten Tick mit demselben Element weiter.
  • Keine Endlosschleife: Wartest du auf Platz, tust du das über die Ticks, nicht in einer Schleife.

Der Check lässt den Erzeuger in drei Szenarien laufen (begrenzte Queue, mal ein langsamer, mal ein schneller, mal ein zeitweise stehender Verbraucher) und prüft: kein Element geht verloren oder wird doppelt abgegeben, die Queue überschreitet ihr Limit nie, höchstens ein erzeugtes Element liegt noch beim Erzeuger, und der Verbraucher wird nicht unnötig ausgebremst. Probiere selbst mit lauf(Erzeuger, 5, lambda t: 2, 100).

Was passiert mit einem Element, das put abgelehnt hat, wenn der Erzeuger im nächsten Tick wieder drankommt? Woher weiß er dann, ob er etwas Neues erzeugen darf? Spiele es mit lauf(Erzeuger, 5, lambda t: 2, 100) durch und zähle, wie viele Elemente fehlen.

class Erzeuger:
    def __init__(self, quelle):
        self.quelle = quelle
        self.wartend = None

    def schritt(self, q):
        for _ in range(3):
            if self.wartend is None:
                self.wartend = self.quelle()
            if not q.put(self.wartend):
                return              # voll: Element behalten, nichts Neues erzeugen
            self.wartend = None

Erzeuger

Das nicht abgegebene Element bleibt in self.wartend, und solange es dort liegt, ruft der Erzeuger quelle() nicht erneut auf. So wandert der Druck vom Verbraucher bis zur Quelle zurück: Der Erzeuger produziert im Takt des Verbrauchers, nichts geht verloren, nichts staut sich außerhalb der Queue.

Übung 3: Timeout und Abbruch durch eine Kette reparieren (ca. 15 Min.)

Die Kette handler, service, datenbank ist geladen (datenbank(dauer) wartet dauer Sekunden). Jede Schicht trägt sich beim Start in offen ein und beim Ende wieder aus (im finally), außerdem schreiben sie nach dem Aufräumen in log. Zusätzlich ist laufe aus Schritt 1 geladen, damit der Check mit virtueller Uhr läuft.

Der folgende handler hat ein Timeout, leakt aber: Nach dem Timeout laufen service und datenbank weiter. Repariere ihn:

  • Dauert der Aufruf kürzer als timeout, gibt handler das Ergebnis von service(dauer) zurück.
  • Sonst gibt er "timeout" zurück, und zwar nach etwa timeout Sekunden, und zu diesem Zeitpunkt ist die Kette bereits aufgeräumt (offen leer).
  • Bricht der Aufrufer handler ab (task.cancel()), muss die Abbruch-Ausnahme durchlaufen, und die Kette darunter wird ebenfalls aufgeräumt.

Der Timeout muss ein Teil der Struktur sein, kein Nebenher-Task: Alles, was unter dem Handler läuft, gehört ihm und muss enden, bevor er zurückkehrt. Welches Werkzeug aus Schritt 3 legt eine Deadline um einen ganzen Block? Und beachte, dass die Abbruch-Ausnahme von außen nicht in deinem except hängen bleiben darf.

async def handler(dauer, timeout):
    try:
        async with asyncio.timeout(timeout):
            return await service(dauer)
    except TimeoutError:
        return "timeout"

handler

asyncio.timeout ist eine Deadline für den ganzen async with-Block. Läuft sie ab, wird die Arbeit im Block abgebrochen (Ausnahme am innersten await, jedes finally räumt auf), und erst danach wird sie zu TimeoutError. Weil service direkt in handler abgewartet wird und kein losgelöster Task ist, gilt dasselbe für einen Abbruch von außen. asyncio.wait_for(service(dauer), timeout) wäre ebenfalls richtig.

Merksatz und Prüfstein

Merksatz: Jede Aufgabe, die du startest, muss enden, bevor ihr Besitzer endet, und ein langsamer Verbraucher muss den Erzeuger bremsen (Backpressure), statt dass Daten still verloren gehen oder der Speicher voll läuft.

Prüfstein: Wie führst du Cancellation und Timeouts sauber durch eine Kette asynchroner Aufrufe?

Ein Timeout ist eine Deadline für die ganze Kette (asyncio.timeout, in JS AbortSignal.timeout plus das Signal in jede Schicht reichen). Beim Ablauf wird die Arbeit von innen nach außen abgebrochen: Die Abbruch-Ausnahme (CancelledError) läuft durch, jede Schicht räumt im finally auf, und der Besitzer kehrt erst danach zurück. Dazu: CancelledError nie verschlucken, Tasks nicht losgelöst starten (TaskGroup statt create_task von Hand), cancel() ist nur eine Bitte und braucht ein Abwarten, und Fehler kommen beim Besitzer an.


Quelle: quellen/konzeptuebersicht-software-grundlagen.md, Abschnitt “3. Nebenläufigkeit” (Event Loop, Backpressure); quellen/konzeptuebersicht-software-fortgeschritten.docx, Schicht 3 (Strukturierte Nebenläufigkeit, Scheduling, Lastverhalten). Die Quellen sind Stichwortlisten.

Über die Quelle hinaus (allgemeines Fachwissen): die virtuelle Uhr als Simulation, die Aussagen zu blockierenden Aufrufen und asyncio.to_thread, die Aussage zu Node-Streams (write() liefert false, drain), asyncio.timeout, TaskGroup und ExceptionGroup, die Abbruch-Regeln. Zahlen in den Ausgaben stammen aus ausgeführtem Code (Python 3.13.9, Node 20).