Files
Orchestrator/bahn/teamlandkarte-mcp/.kiro/specs/llm-batch-matching/design.md
T
ankn a5f8fb49ab Migrate all repos into monorepo context folders
Bahn: aisupport, Analyse-O2C-C2S, awesome-bahn-mcp-servers, beam-mcp,
      Confluence_Bot, db-planet-mcp-server, O2C-Harness, project-audit,
      Projekt-KIQ-HP, teamlandkarte-mcp
Dhive: Jury-Voting
Privat: CV, NoteGraph (NOTE: NoteGraph needs complete redo after consolidation)
Shared: AI-Orchestrator, OrgMyLife, power_skills_and_more
Shared/references: symphony (read-only)

Bahn repos remain available as independent remotes - this monorepo
pulls them in via subtree, the originals are untouched.
2026-06-30 20:39:52 +02:00

13 KiB
Raw Blame History

Design: LLM Batch-Matching (Concurrent Requests)

Übersicht

Das bestehende LlmFulltextMatcher-Modul führt LLM-Aufrufe sequenziell aus jeder Kandidat wartet auf die Antwort des vorherigen. Bei 30 Kandidaten mit je ~2s Latenz ergibt das ~60s Gesamtlaufzeit.

Die Lösung ersetzt die sequenzielle for-Schleife durch asyncio.gather() mit einem asyncio.Semaphore zur Begrenzung der Parallelität. Die öffentliche Schnittstelle (match_capacities, match_tasks) bleibt unverändert. Die Concurrency wird über config.toml konfigurierbar gemacht.

Erwarteter Effekt: Bei max_concurrency=5 und 30 Kandidaten sinkt die Laufzeit von ~60s auf ~12s (6 Batches × 2s statt 30 × 2s).

Architektur

graph TD
    A[MCP Server / Aufrufer] -->|match_capacities / match_tasks| B[LlmFulltextMatcher]
    B -->|asyncio.gather + Semaphore| C[_categorize_one Task 1]
    B -->|asyncio.gather + Semaphore| D[_categorize_one Task 2]
    B -->|asyncio.gather + Semaphore| E[_categorize_one Task N]
    C --> F[AzureOpenAIClient.chat_completion]
    D --> F
    E --> F
    F --> G[Azure OpenAI API]

Die Architektur bleibt flach: Es wird kein neues Modul oder neue Klasse eingeführt. Die Änderung betrifft ausschließlich die interne Ablaufsteuerung in LlmFulltextMatcher und die Konfigurationsschicht.

Designentscheidungen

  1. asyncio.Semaphore statt Thread-Pool: Das Projekt ist bereits vollständig async (FastMCP, AsyncAzureOpenAI). Ein Semaphore ist der idiomatische Mechanismus zur Begrenzung von I/O-Concurrency in asyncio.

  2. Kein separates Batch-Modul: Die Änderung ist minimal und lokal. Ein eigenes batch_executor.py wäre Over-Engineering für eine ~20-Zeilen-Änderung.

  3. Semaphore im Matcher, nicht im Client: Der AzureOpenAIClient bleibt unverändert. Die Concurrency-Steuerung liegt beim Aufrufer (Matcher), da verschiedene Aufrufer unterschiedliche Limits haben könnten.

  4. Validierung bei Config-Load: Ungültige max_concurrency-Werte (< 1 oder > 20) werden beim Start abgefangen, nicht erst beim ersten Matching-Aufruf.

Komponenten und Schnittstellen

1. AzureOpenAIConfig (config.py)

Neues Feld:

@dataclass
class AzureOpenAIConfig:
    # ... bestehende Felder ...
    max_concurrency: int = 5

2. _parse_azure_openai (config.py)

Erweiterte Parsing-Logik:

def _parse_azure_openai(cfg: dict) -> AzureOpenAIConfig:
    max_concurrency = int(cfg.get("max_concurrency", 5))
    if max_concurrency < 1:
        raise ConfigError(
            "azure_openai.max_concurrency must be >= 1, "
            f"got {max_concurrency}"
        )
    if max_concurrency > 20:
        raise ConfigError(
            "azure_openai.max_concurrency must be <= 20, "
            f"got {max_concurrency}"
        )
    return AzureOpenAIConfig(
        # ... bestehende Felder ...
        max_concurrency=max_concurrency,
    )

3. LlmFulltextMatcher (matching/llm_fulltext_matcher.py)

Geänderte Konstruktor-Signatur:

class LlmFulltextMatcher:
    def __init__(
        self,
        *,
        db: DBClient,
        client: AzureOpenAIClient,
        rationale_max_chars: int = 280,
        max_concurrency: int = 5,  # NEU
    ) -> None:
        self._db = db
        self._client = client
        self._rationale_max_chars = rationale_max_chars
        self._semaphore = asyncio.Semaphore(max_concurrency)
        self._max_concurrency = max_concurrency

Neue interne Methode:

async def _categorize_one_throttled(
    self,
    *,
    item_id: str,
    user_prompt: str,
    raw: dict,
) -> tuple[LlmFulltextItem | None, LlmFulltextError | None]:
    """Wrapper um _categorize_one mit Semaphore-Begrenzung."""
    async with self._semaphore:
        return await self._categorize_one(
            item_id=item_id,
            user_prompt=user_prompt,
            raw=raw,
        )

Geänderte match_capacities / match_tasks (Kernänderung):

async def match_capacities(self, *, task_profile, capacities) -> LlmFulltextResult:
    # ... Profil-Aufbau wie bisher ...

    LOGGER.info(
        "Batch-Matching gestartet: %d Kandidaten, max_concurrency=%d",
        len(capacities), self._max_concurrency,
    )
    start_time = time.monotonic()

    tasks = [
        self._categorize_one_throttled(
            item_id=cap_id,
            user_prompt=user_prompt,
            raw=raw,
        )
        for cap_id, user_prompt, raw in prepared
    ]
    results = await asyncio.gather(*tasks)

    elapsed = time.monotonic() - start_time

    # Ergebnisse zuordnen
    for item, error in results:
        if item is not None:
            by_category[item.category].append(item)
        elif error is not None:
            errors.append(error)

    LOGGER.info(
        "Batch-Matching abgeschlossen: %.1fs, %d kategorisiert, %d Fehler",
        elapsed, sum(len(v) for v in by_category.values()), len(errors),
    )
    # ... Sortierung wie bisher ...

4. MCP Server (mcp_server.py)

Übergabe des neuen Parameters bei Instanziierung:

llm_fulltext_matcher = LlmFulltextMatcher(
    db=db_client,
    client=azure_client,
    max_concurrency=cfg.azure_openai.max_concurrency,  # NEU
)

5. config.toml

Neuer optionaler Schlüssel:

[azure_openai]
# ... bestehende Schlüssel ...

# Maximale Anzahl paralleler LLM-Anfragen (1-20, Standard: 5).
# Höhere Werte beschleunigen das Matching, können aber Rate-Limits auslösen.
# max_concurrency = 5

Datenmodelle

Keine neuen Datenmodelle erforderlich. Die bestehenden Strukturen bleiben unverändert:

  • LlmFulltextResult (Rückgabetyp) unverändert
  • LlmFulltextItem unverändert
  • LlmFulltextError unverändert
  • AzureOpenAIConfig erweitert um max_concurrency: int = 5

Die Erweiterung von AzureOpenAIConfig ist abwärtskompatibel (Standardwert vorhanden).

Correctness Properties

Eine Property ist eine Eigenschaft oder ein Verhalten, das über alle gültigen Ausführungen eines Systems hinweg gelten muss im Wesentlichen eine formale Aussage darüber, was das System tun soll. Properties bilden die Brücke zwischen menschenlesbaren Spezifikationen und maschinell verifizierbaren Korrektheitsgarantien.

Property 1: Vollständigkeit der Ergebnisse (Partition)

Für jede Liste von Kandidaten (Capacities oder Tasks) und jede Konfiguration von max_concurrency, muss die Summe aller Items in by_category plus die Anzahl der Einträge in errors exakt der Anzahl der Eingabe-Kandidaten entsprechen. Kein Kandidat darf verloren gehen oder doppelt erscheinen.

Validates: Requirements 1.1, 1.4, 3.2, 4.2, 4.3

Property 2: Concurrency-Begrenzung

Für jede Anzahl von Kandidaten und jeden gültigen max_concurrency-Wert, darf zu keinem Zeitpunkt die Anzahl gleichzeitig laufender LLM-Aufrufe den konfigurierten max_concurrency-Wert überschreiten.

Validates: Requirements 1.2

Property 3: Deterministische Sortierung

Für jede Liste von Kandidaten und jede Zuordnung von Kategorien, muss das Ergebnis innerhalb jeder Kategorie aufsteigend nach item_id (lexikographisch) sortiert sein, und die Fehlerliste muss ebenfalls nach item_id sortiert sein.

Validates: Requirements 3.1

Property 4: Eingabereihenfolge-Unabhängigkeit

Für jede Permutation der Eingabe-Kandidatenliste muss das Ergebnis (by_category und errors) identisch sein die Reihenfolge der Eingabe hat keinen Einfluss auf die Ausgabe.

Validates: Requirements 3.3

Property 5: Korrekte Zuordnung (Response-Mapping)

Für jeden Kandidaten in der Eingabeliste muss die zugehörige LLM-Antwort exakt dem richtigen Kandidaten zugeordnet werden. Wenn der Mock für Kandidat X die Kategorie "Top" zurückgibt, muss das Item mit item_id=X in by_category["Top"] erscheinen.

Validates: Requirements 3.4

Property 6: Config-Validierung

Für jeden Integer-Wert n: Das Parsen von azure_openai.max_concurrency = n muss genau dann erfolgreich sein, wenn 1 <= n <= 20. Für n < 1 oder n > 20 muss ein ConfigError geworfen werden.

Validates: Requirements 2.1, 2.3, 2.4

Property 7: Logging-Konsistenz

Für jede Ausführung von match_capacities oder match_tasks mit mindestens einem Kandidaten müssen die INFO-Log-Nachrichten (Start und Ende) die korrekte Kandidatenanzahl, das konfigurierte Concurrency-Limit, die Anzahl erfolgreicher Kategorisierungen und die Anzahl der Fehler enthalten. Die Summe von Erfolgen und Fehlern im Log muss der Eingabeanzahl entsprechen.

Validates: Requirements 6.1, 6.2

Fehlerbehandlung

Fehlerszenario Verhalten Auswirkung auf andere Kandidaten
Einzelner LLM-Aufruf schlägt fehl (nach Retries) Kandidat wird in errors-Liste aufgenommen Keine andere Kandidaten laufen unabhängig weiter
Alle LLM-Aufrufe schlagen fehl Leere by_category, vollständige errors-Liste Kein Exception-Wurf, normaler Return
max_concurrency ungültig (< 1 oder > 20) ConfigError beim Server-Start Server startet nicht
asyncio.TimeoutError in einem Aufruf Wird vom bestehenden Retry-Mechanismus in AzureOpenAIClient behandelt Keine
HTTP 429 (Rate Limit) Exponentielles Backoff im AzureOpenAIClient (bestehendes Verhalten) Keine direkte; Semaphore hält Slot belegt bis Retry abgeschlossen

Fehler-Isolation

Die Verwendung von asyncio.gather(*tasks) (ohne return_exceptions=True) in Kombination mit der bestehenden try/except-Logik in _categorize_one stellt sicher, dass:

  • Jeder Task seine eigenen Exceptions fängt und als LlmFulltextError zurückgibt
  • Kein einzelner Fehler die gesamte gather-Operation abbricht
  • Die Semaphore auch im Fehlerfall korrekt freigegeben wird (async context manager)

Teststrategie

Property-Based Tests (pytest + hypothesis)

Die Property-Tests verwenden die Bibliothek hypothesis (bereits im Python-Ökosystem etabliert). Jeder Test wird mit mindestens 100 Iterationen konfiguriert.

Property Testansatz Generator
P1: Vollständigkeit Generiere zufällige Kandidatenlisten mit zufälligem Mix aus Erfolg/Fehler-Mocks st.lists(st.builds(Capacity, ...))
P2: Concurrency-Begrenzung Mock-Client mit Counter für gleichzeitige Aufrufe; prüfe max_concurrent <= max_concurrency st.integers(min_value=1, max_value=20) für concurrency
P3: Deterministische Sortierung Generiere Ergebnisse mit zufälligen Kategorien; prüfe Sortierung st.lists(st.sampled_from(categories))
P4: Eingabereihenfolge-Unabhängigkeit Generiere Liste, permutiere, vergleiche Ergebnisse st.permutations(candidates)
P5: Korrekte Zuordnung Mock gibt item_id-spezifische Kategorien zurück; prüfe Mapping st.dictionaries(st.text(), st.sampled_from(categories))
P6: Config-Validierung Generiere Integers im Bereich [-100, 100]; prüfe Erfolg/Fehler st.integers(min_value=-100, max_value=100)
P7: Logging-Konsistenz Capture Logs; prüfe Zahlen gegen tatsächliche Ergebnisse st.lists(st.builds(Capacity, ...))

Jeder Test wird mit einem Kommentar getaggt:

# Feature: llm-batch-matching, Property 1: Vollständigkeit der Ergebnisse

Unit Tests

Unit Tests ergänzen die Property-Tests für spezifische Szenarien:

  • Leere Eingabe: match_capacities(capacities=[]) gibt sofort leeres Ergebnis zurück
  • Alle Fehler: Wenn alle LLM-Aufrufe fehlschlagen, keine Exception, vollständige Fehlerliste
  • Default-Wert: LlmFulltextMatcher() ohne max_concurrency verwendet 5
  • Config-Parsing: Fehlender max_concurrency-Schlüssel ergibt Standardwert 5
  • Integration: End-to-End-Test mit gemocktem AzureOpenAIClient und 10 Kandidaten

Testinfrastruktur

  • Mock-Client: Ein FakeAzureOpenAIClient der konfigurierbare Antworten (Erfolg/Fehler/Delay) pro item_id liefert
  • Concurrency-Tracker: Ein Wrapper der die maximale Anzahl gleichzeitiger Aufrufe misst (via asyncio.Lock + Counter)
  • Log-Capture: pytest caplog Fixture für Log-Assertions