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

Verwandte Artikel

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.


← Zurück zum Dokumentationsindex