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.
288 lines
13 KiB
Markdown
288 lines
13 KiB
Markdown
# 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
|
||
|
||
```mermaid
|
||
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:
|
||
|
||
```python
|
||
@dataclass
|
||
class AzureOpenAIConfig:
|
||
# ... bestehende Felder ...
|
||
max_concurrency: int = 5
|
||
```
|
||
|
||
### 2. `_parse_azure_openai` (config.py)
|
||
|
||
Erweiterte Parsing-Logik:
|
||
|
||
```python
|
||
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:
|
||
|
||
```python
|
||
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:
|
||
|
||
```python
|
||
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):
|
||
|
||
```python
|
||
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:
|
||
|
||
```python
|
||
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:
|
||
|
||
```toml
|
||
[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:
|
||
|
||
```python
|
||
# 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
|