Elaborazione Incrementale: Aggiornamenti Rapidi della Pipeline

Questo articolo spiega come la pipeline SEO utilizza l'elaborazione incrementale per essere eseguita in secondi invece che in ore.

Il Problema: La Riprocessazione Completa è Lenta

Eseguire l'intera pipeline da zero richiede ore:

  • Step 0 (Embedding delle sorgenti): 15 minuti (65,000 prodotti)

  • Step 1 (Recupero query): 10 minuti (chiamate API)

  • Step 2 (Clustering query): 30 minuti (similarità 65K×65K)

  • Step 3 (Mappature frasi): 20 minuti (embedding + matching)

  • Step 4 (Matching prodotti): 45 minuti (query × prodotti)

  • Step 5 (Ricerche correlate): 25 minuti (similarità query × query)

Totale: ~2.5 ore per la pipeline completa

Problema: Gli aggiornamenti giornalieri sprecherebbero 2.5 ore ricalcolando dati invariati.

La Soluzione: Strategia Incrementale a Tre Livelli

Utilizziamo tre tecniche per saltare il lavoro non necessario:

1. Salto dello Step (Granularità Grossolana)

Salta interi step se l'output è aggiornato e lo script non è cambiato.

2. Embedding Incrementale (Granularità Media)

Esegui l'embedding solo degli elementi nuovi/modificati, riutilizza gli embedding in cache.

3. Checkpointing (Granularità Fine)

Salva i progressi durante le operazioni lunghe, riprendi dal checkpoint in caso di fallimento.

Strategia di Salto dello Step

Come Funziona

Prima di ogni step, controlla:

L'output esiste? Se no, esegui lo step.

Età dell'output: Se più vecchio di 7 giorni, esegui lo step.

Lo script è cambiato? Se lo script è stato modificato dopo la generazione dell'output, esegui lo step.

Tutti i controlli superati? Salta lo step.

Implementazione

def should_skip_step(output_path, script_path, days=7):
    # Controlla se l'output esiste
    if not os.path.exists(output_path):
# ... (dettagli implementativi omessi)

Utilizzo

Ogni script controlla all'avvio:

from seo_common import should_skip_step

if should_skip_step(SEO_SOURCE_EMBEDDINGS_PATH, __file__):
    print("✓ Salto: L'output è aggiornato e lo script non è cambiato")
    return

Vantaggi

Esecuzioni giornaliere veloci: La maggior parte degli step viene saltata se i dati non cambiano

Invalidazione automatica: Le modifiche allo script innescano una riesecuzione

Freschezza configurabile: Regola il parametro days per ogni step

Strategia di Embedding Incrementale

Come Funziona

Quando si esegue l'embedding di elementi (prodotti, query, frasi):

Carica la cache: Leggi gli elementi precedentemente embedded e le loro chiavi

Confronta le chiavi: Identifica elementi nuovi, modificati ed eliminati

Esegui embedding solo dei nuovi: Esegui l'embedding solo degli elementi non presenti in cache

Unisci: Combina gli embedding in cache con quelli nuovi nell'ordine corretto

Salva: Scrivi la cache aggiornata

Implementazione

La funzione incremental_embed_with_keys gestisce questo processo:

def incremental_embed_with_keys(
    items,           # Elementi correnti da elaborare
    keys,            # Chiavi uniche per gli elementi
# ... (dettagli implementativi omessi)

Tassi di Cache Hit

Tassi tipici di cache hit dopo la prima esecuzione:

Dati sorgente (prodotti, parti, articoli):

  • Prima esecuzione: 0% (embedding di tutti i 65,000 elementi)

  • Esecuzione giornaliera: 99.5% (solo ~300 elementi nuovi/modificati)

Query (da GSC, Ads, live):

  • Prima esecuzione: 0% (embedding di tutte le 65,000 query)

  • Esecuzione giornaliera: 99.2% (solo ~500 nuove query)

Mappature frasi:

  • Prima esecuzione: 0% (embedding di tutte le 5,000 frasi)

  • Esecuzione giornaliera: 99.8% (solo ~10 nuove frasi)

Impatto sulle Prestazioni

Prima esecuzione (cache fredda):

  • Embedding sorgenti: 15 minuti (65,000 elementi)

  • Embedding query: 10 minuti (65,000 query)

  • Embedding frasi: 2 minuti (5,000 frasi)

Esecuzione giornaliera (cache calda):

  • Embedding sorgenti: 10 secondi (300 elementi, tasso di cache hit 99.5%)

  • Embedding query: 5 secondi (500 query, tasso di cache hit 99.2%)

  • Embedding frasi: 1 secondo (10 frasi, tasso di cache hit 99.8%)

Incremento di velocità: 90-180× più veloce

Strategia di Checkpointing

Come Funziona

Per operazioni a lunga esecuzione (embedding di 65,000 elementi):

Elaborazione a lotti: Elabora gli elementi in lotti (es. 1,000 elementi)

Salva checkpoint: Dopo ogni lotto, salva i risultati accumulati

Riprendi in caso di fallimento: Se il processo si interrompe, riprendi dall'ultimo checkpoint

Salvataggio finale: Dopo tutti i lotti, salva i risultati completi

Implementazione

Il checkpointing è integrato in incremental_embed_with_keys:

checkpoint_every = 1000  # Salva ogni 1,000 elementi

embeddings_list = []
# ... (dettagli implementativi omessi)

Vantaggi

Recupero da interruzione: Riprendi dall'ultimo checkpoint invece di ricominciare da zero

Visibilità dei progressi: Vedi i progressi ogni 1,000 elementi

Efficienza della memoria: Elabora in lotti, non caricare tutto in una volta

Integrazione nella Pipeline

L'elaborazione incrementale è utilizzata in più step:

Step 0: Embedding Dati Sorgente

Incrementale: Esegui l'embedding solo di prodotti, parti, articoli nuovi/modificati

Checkpointing: Salva ogni 1,000 elementi

Logica di salto: Salta se l'output ha meno di 7 giorni e lo script non è cambiato

Vedi: Embedding Dati Sorgente

Step 1: Recupero Query

Incrementale: Le chiamate API recuperano solo dati nuovi (dall'ultima esecuzione)

Logica di salto: Salta se l'output ha meno di 1 giorno

Vedi: Recupero Query

Step 3b: Embedding Query

Incrementale: Esegui l'embedding solo delle nuove query

Checkpointing: Salva ogni 1,000 query

Logica di salto: Salta se l'output ha meno di 7 giorni e lo script non è cambiato

Vedi: Embedding Query

Step 4: Espansione Mappature Frasi

Incrementale: Esegui l'embedding solo delle nuove frasi

Checkpointing: Salva ogni 1,000 frasi

Logica di salto: Salta se l'output ha meno di 7 giorni e lo script non è cambiato

Vedi: Mappature Frasi-Filtri

Step 6: Matching Prodotti

Incrementale: Abbina solo le nuove query

Logica di salto: Salta se l'output ha meno di 7 giorni e lo script non è cambiato

Vedi: Matching Prodotti

Configurazione

L'elaborazione incrementale è configurata per ogni step:

Soglia di Freschezza

# Salta se l'output ha meno di 7 giorni (default)
should_skip_step(output_path, script_path, days=7)

# Salta se l'output ha meno di 1 giorno (per dati che cambiano frequentemente)
should_skip_step(output_path, script_path, days=1)

Frequenza Checkpoint

# Salva ogni 1,000 elementi (default)
incremental_embed_with_keys(..., checkpoint_every=1000)

# Salva ogni 5,000 elementi (per elaborazione più veloce, meno sicurezza)
incremental_embed_with_keys(..., checkpoint_every=5000)

Dimensione Lotto

# Embedding di 32 elementi per lotto (default, bilanciato)
incremental_embed_with_keys(..., batch_size=32)

# Embedding di 64 elementi per lotto (più veloce su GPU, più memoria)
incremental_embed_with_keys(..., batch_size=64)

Monitoraggio e Debug

Statistiche Cache

Ogni step stampa le statistiche della cache:

✓ Cache esistente trovata, controllo delle modifiche...
  Esistenti: 65,000 elementi
  Correnti:  65,300 elementi
  Riutilizzo: 64,800 embedding
  Nuovi:     500 elementi da elaborare

Messaggi di Salto

Quando gli step vengono saltati:

✓ Salto di 0_embed_source_data.py: L'output è aggiornato e lo script non è cambiato.

Messaggi Checkpoint

Durante operazioni lunghe:

Embedding di 65,000 elementi (checkpoint ogni 1,000)...
  Lotto 0-1000...
    ✓ Checkpoint salvato (1,000 totali)
  Lotto 1000-2000...
    ✓ Checkpoint salvato (2,000 totali)
  ...

Riferimenti

Concetti Tecnici

Articoli Correlati

Riepilogo

L'elaborazione incrementale rende la pipeline 90-180× più veloce:

Strategia a tre livelli:

  • ✅ Salto dello step (salta interi step se l'output è aggiornato)

  • ✅ Embedding incrementale (elabora solo elementi nuovi/modificati)

  • ✅ Checkpointing (salva i progressi, riprendi in caso di fallimento)

Prestazioni:

  • ✅ Prima esecuzione: ~2.5 ore (pipeline completa)

  • ✅ Esecuzione giornaliera: ~5 minuti (aggiornamenti incrementali)

  • ✅ Tassi di cache hit: 99%+ dopo la prima esecuzione

Vantaggi:

  • ✅ Aggiornamenti giornalieri veloci (minuti invece di ore)

  • ✅ Invalidazione automatica (le modifiche allo script innescano una riesecuzione)

  • ✅ Recupero da interruzione (riprendi dal checkpoint)

  • ✅ Efficienza della memoria (elaborazione a lotti)

Questa strategia permette esecuzioni giornaliere della pipeline senza sprecare risorse di calcolo su dati invariati.


← Torna all'Indice della Documentazione