Inkrementelle Verarbeitung: Schnelle Pipeline-Aktualisierungen
Dieser Artikel erklärt, wie die SEO-Pipeline inkrementelle Verarbeitung nutzt, um in Sekunden statt Stunden zu laufen.
Das Problem: Vollständige Neuverarbeitung ist langsam
Die gesamte Pipeline von Grund auf neu auszuführen, dauert Stunden:
-
Schritt 0 (Source Embedding): 15 Minuten (65.000 Produkte)
-
Schritt 1 (Query Fetching): 10 Minuten (API-Aufrufe)
-
Schritt 2 (Query Clustering): 30 Minuten (65K×65K Ähnlichkeit)
-
Schritt 3 (Phrase Mappings): 20 Minuten (Embedding + Matching)
-
Schritt 4 (Product Matching): 45 Minuten (Queries × Produkte)
-
Schritt 5 (Related Searches): 25 Minuten (Query × Query Ähnlichkeit)
Gesamt: ~2,5 Stunden für die vollständige Pipeline
Problem: Tägliche Updates würden 2,5 Stunden mit der Neuberechnung unveränderter Daten verschwenden.
Die Lösung: Drei-Schichten-Inkrementell-Strategie
Wir verwenden drei Techniken, um unnötige Arbeit zu überspringen:
1. Schrittüberspringen (Grobkörnig)
Überspringe ganze Schritte, wenn die Ausgabe aktuell ist und das Skript unverändert.
2. Inkrementelles Embedding (Mittelkörnig)
Embedde nur neue/geänderte Elemente, verwende zwischengespeicherte Embeddings wieder.
3. Checkpointing (Feinkörnig)
Speichere den Fortschritt während langer Operationen, setze bei einem Fehler vom Checkpoint fort.
Schrittüberspringen-Strategie
So funktioniert es
Vor jedem Schritt prüfen:
Existiert die Ausgabe? Wenn nein, führe den Schritt aus.
Ausgabealter: Wenn älter als 7 Tage, führe den Schritt aus.
Skript geändert? Wenn das Skript seit der Erstellung der Ausgabe geändert wurde, führe den Schritt aus.
Alle Prüfungen bestanden? Überspringe den Schritt.
Implementierung
def should_skip_step(output_path, script_path, days=7):
# Check if output exists
if not os.path.exists(output_path):
# ... (implementation details omitted)
Verwendung
Jedes Skript prüft beim Start:
from seo_common import should_skip_step
if should_skip_step(SEO_SOURCE_EMBEDDINGS_PATH, __file__):
print("✓ Skipping: Output is fresh and script unchanged")
return
Vorteile
Schnelle tägliche Läufe: Die meisten Schritte werden übersprungen, wenn Daten unverändert sind
Automatische Ungültigmachung: Skriptänderungen lösen einen Neulauf aus
Konfigurierbare Aktualität: Passe den days-Parameter pro Schritt an
Inkrementelle Embedding-Strategie
So funktioniert es
Beim Embedding von Elementen (Produkte, Queries, Phrasen):
Cache laden: Lese zuvor embeddede Elemente und ihre Schlüssel
Schlüssel vergleichen: Identifiziere neue, geänderte und gelöschte Elemente
Nur neue embedden: Embedde nur Elemente, die nicht im Cache sind
Zusammenführen: Kombiniere zwischengespeicherte Embeddings mit neuen Embeddings in der richtigen Reihenfolge
Speichern: Schreibe den aktualisierten Cache
Implementierung
Die Funktion incremental_embed_with_keys behandelt dies:
def incremental_embed_with_keys(
items, # Current items to embed
keys, # Unique keys for items
# ... (implementation details omitted)
Cache-Trefferquoten
Typische Cache-Trefferquoten nach dem ersten Lauf:
Quelldaten (Produkte, Teile, Artikel):
-
Erster Lauf: 0% (alle 65.000 Elemente embedden)
-
Täglicher Lauf: 99,5% (nur ~300 neue/geänderte Elemente)
Queries (von GSC, Ads, live):
-
Erster Lauf: 0% (alle 65.000 Queries embedden)
-
Täglicher Lauf: 99,2% (nur ~500 neue Queries)
Phrasen-Mappings:
-
Erster Lauf: 0% (alle 5.000 Phrasen embedden)
-
Täglicher Lauf: 99,8% (nur ~10 neue Phrasen)
Leistungsauswirkung
Erster Lauf (kalter Cache):
-
Source Embedding: 15 Minuten (65.000 Elemente)
-
Query Embedding: 10 Minuten (65.000 Queries)
-
Phrase Embedding: 2 Minuten (5.000 Phrasen)
Täglicher Lauf (warmer Cache):
-
Source Embedding: 10 Sekunden (300 Elemente, 99,5% Trefferquote)
-
Query Embedding: 5 Sekunden (500 Queries, 99,2% Trefferquote)
-
Phrase Embedding: 1 Sekunde (10 Phrasen, 99,8% Trefferquote)
Beschleunigung: 90-180× schneller
Checkpointing-Strategie
So funktioniert es
Für langlaufende Operationen (Embedding von 65.000 Elementen):
Batch-Verarbeitung: Verarbeite Elemente in Batches (z.B. 1.000 Elemente)
Checkpoint speichern: Nach jedem Batch speichere die akkumulierten Ergebnisse
Bei Fehler fortsetzen: Wenn der Prozess abstürzt, setze vom letzten Checkpoint fort
Endgültig speichern: Nach allen Batches speichere die vollständigen Ergebnisse
Implementierung
Checkpointing ist in incremental_embed_with_keys integriert:
checkpoint_every = 1000 # Save every 1,000 items
embeddings_list = []
# ... (implementation details omitted)
Vorteile
Absturzwiederherstellung: Fortsetzen vom letzten Checkpoint statt von vorne
Fortschrittssichtbarkeit: Siehe Fortschritt alle 1.000 Elemente
Speichereffizienz: Verarbeite in Batches, lade nicht alles auf einmal
Integration über die Pipeline hinweg
Inkrementelle Verarbeitung wird in mehreren Schritten verwendet:
Schritt 0: Quelldaten-Embedding
Inkrementell: Nur neue/geänderte Produkte, Teile, Artikel embedden
Checkpointing: Speichere alle 1.000 Elemente
Überspring-Logik: Überspringe, wenn Ausgabe < 7 Tage alt und Skript unverändert
Siehe: Quelldaten-Embedding
Schritt 1: Query Fetching
Inkrementell: API-Aufrufe holen nur neue Daten (seit dem letzten Lauf)
Überspring-Logik: Überspringe, wenn Ausgabe < 1 Tag alt
Siehe: Query Fetching
Schritt 3b: Query Embedding
Inkrementell: Nur neue Queries embedden
Checkpointing: Speichere alle 1.000 Queries
Überspring-Logik: Überspringe, wenn Ausgabe < 7 Tage alt und Skript unverändert
Siehe: Query Embedding
Schritt 4: Phrasen-Mapping-Erweiterung
Inkrementell: Nur neue Phrasen embedden
Checkpointing: Speichere alle 1.000 Phrasen
Überspring-Logik: Überspringe, wenn Ausgabe < 7 Tage alt und Skript unverändert
Siehe: Phrasen-zu-Filter-Mappings
Schritt 6: Product Matching
Inkrementell: Nur neue Queries matchen
Überspring-Logik: Überspringe, wenn Ausgabe < 7 Tage alt und Skript unverändert
Siehe: Product Matching
Konfiguration
Inkrementelle Verarbeitung wird pro Schritt konfiguriert:
Aktualitätsschwelle
# Skip if output < 7 days old (default)
should_skip_step(output_path, script_path, days=7)
# Skip if output < 1 day old (for frequently changing data)
should_skip_step(output_path, script_path, days=1)
Checkpoint-Häufigkeit
# Save every 1,000 items (default)
incremental_embed_with_keys(..., checkpoint_every=1000)
# Save every 5,000 items (for faster processing, less safety)
incremental_embed_with_keys(..., checkpoint_every=5000)
Batch-Größe
# Embed 32 items per batch (default, balanced)
incremental_embed_with_keys(..., batch_size=32)
# Embed 64 items per batch (faster on GPU, more memory)
incremental_embed_with_keys(..., batch_size=64)
Überwachung und Fehlerbehebung
Cache-Statistiken
Jeder Schritt gibt Cache-Statistiken aus:
✓ Found existing cache, checking for changes...
Existing: 65,000 items
Current: 65,300 items
Reusing: 64,800 embeddings
New: 500 items to embed
Überspring-Nachrichten
Wenn Schritte übersprungen werden:
✓ Skipping 0_embed_source_data.py: Output is fresh and script unchanged.
Checkpoint-Nachrichten
Während langer Operationen:
Embedding 65,000 items (checkpointing every 1,000)...
Batch 0-1000...
✓ Checkpoint saved (1,000 total)
Batch 1000-2000...
✓ Checkpoint saved (2,000 total)
...
Referenzen
Technische Konzepte
-
Inkrementelles Lernen - Wikipedia
-
Checkpointing - Wikipedia
-
Caching - Wikipedia
Verwandte Artikel
-
SEO-Pipeline-Überblick - Komplette Pipeline-Architektur
-
Quelldaten-Embedding - Inkrementelles Produkt-Embedding
-
Query Embedding - Inkrementelles Query-Embedding
-
Phrasen-zu-Filter-Mappings - Inkrementelles Phrasen-Embedding
Zusammenfassung
Inkrementelle Verarbeitung macht die Pipeline 90-180× schneller:
Drei-Schichten-Strategie:
-
✅ Schrittüberspringen (überspringe ganze Schritte, wenn Ausgabe aktuell)
-
✅ Inkrementelles Embedding (nur neue/geänderte Elemente embedden)
-
✅ Checkpointing (speichere Fortschritt, setze bei Fehler fort)
Leistung:
-
✅ Erster Lauf: ~2,5 Stunden (vollständige Pipeline)
-
✅ Täglicher Lauf: ~5 Minuten (inkrementelle Updates)
-
✅ Cache-Trefferquoten: 99%+ nach dem ersten Lauf
Vorteile:
-
✅ Schnelle tägliche Updates (Minuten statt Stunden)
-
✅ Automatische Ungültigmachung (Skriptänderungen lösen Neulauf aus)
-
✅ Absturzwiederherstellung (Fortsetzen vom Checkpoint)
-
✅ Speichereffizient (Batch-Verarbeitung)
Diese Strategie ermöglicht tägliche Pipeline-Läufe, ohne Rechenleistung für unveränderte Daten zu verschwenden.