Python: scegliere tra thread e processi con un caso pratico di pipeline produttore-consumatore

by theArchitect
SHARE
Python: scegliere tra thread e processi con un caso pratico di pipeline produttore-consumatore
© Guida-HTML5.it

Introduzione

Quando si parla di multithreading e multiprocessing in Python, il punto non è solo “come far girare più cose insieme”, ma soprattutto come scegliere l’approccio giusto in base al problema reale. Un errore molto comune è usare thread e processi in modo generico, senza considerare la natura del carico di lavoro.

Un sotto-argomento molto utile, pratico e spesso sottovalutato è la costruzione di una pipeline produttore-consumatore ibrida: i thread raccolgono o preparano i dati, mentre i processi eseguono il lavoro CPU-bound più pesante. Questo modello è perfetto per casi come:

  • lettura di file o stream da rete con I/O frequente;
  • pre-elaborazione leggera dei dati;
  • analisi o trasformazioni intensive su grandi volumi di informazioni;
  • riduzione dei tempi complessivi sfruttando bene CPU e I/O.

In questo tutorial vedremo come progettare una soluzione robusta usando thread per l’I/O e processi per il calcolo, con un esempio concreto e facilmente adattabile a progetti reali.

Codice completo

from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor, as_completed
from pathlib import Path
import os
import time


# -----------------------------
# Fase 1: produttore (I/O)
# -----------------------------
def leggi_file(percorso: Path) -> str:
    """Legge il contenuto di un file di testo."""
    with percorso.open("r", encoding="utf-8") as f:
        return f.read()


def trova_file_testo(cartella: Path):
    """Genera i file .txt presenti in una cartella."""
    for path in cartella.iterdir():
        if path.is_file() and path.suffix == ".txt":
            yield path


# -----------------------------
# Fase 2: lavoro CPU-bound
# -----------------------------
def analizza_testo(testo: str) -> dict:
    """
    Esegue un´analisi volutamente CPU-bound:
    - conta parole
    - conta caratteri
    - calcola una piccola firma numerica
    """
    parole = testo.split()
    numero_parole = len(parole)
    numero_caratteri = len(testo)

    # Simuliamo un´elaborazione più costosa, utile per mostrare il multiprocessing
    firma = 0
    for ch in testo:
        firma = (firma * 31 + ord(ch)) % 1_000_000_007

    return {
        "parole": numero_parole,
        "caratteri": numero_caratteri,
        "firma": firma,
    }


# -----------------------------
# Pipeline ibrida
# -----------------------------
def pipeline(cartella: str):
    cartella_path = Path(cartella)

    if not cartella_path.exists():
        raise FileNotFoundError(f"La cartella {cartella} non esiste.")

    file_testo = list(trova_file_testo(cartella_path))
    if not file_testo:
        print("Nessun file .txt trovato.")
        return

    risultati = []

    # Thread: ottimi per I/O, ad esempio lettura file concorrente
    with ThreadPoolExecutor(max_workers=4) as thread_pool:
        future_lettura = {
            thread_pool.submit(leggi_file, path): path
            for path in file_testo
        }

        contenuti = []
        for future in as_completed(future_lettura):
            path = future_lettura[future]
            try:
                testo = future.result()
                contenuti.append((path.name, testo))
            except Exception as e:
                print(f"Errore nella lettura di {path.name}: {e}")

    # Processi: ottimi per lavoro CPU-bound
    with ProcessPoolExecutor(max_workers=os.cpu_count() or 2) as process_pool:
        future_analisi = {
            process_pool.submit(analizza_testo, testo): nome_file
            for nome_file, testo in contenuti
        }

        for future in as_completed(future_analisi):
            nome_file = future_analisi[future]
            try:
                risultato = future.result()
                risultati.append((nome_file, risultato))
            except Exception as e:
                print(f"Errore nell´analisi di {nome_file}: {e}")

    # Ordinamento finale per leggibilità
    risultati.sort(key=lambda x: x[0])

    print("nRisultati finali:")
    for nome_file, dati in risultati:
        print(
            f"- {nome_file}: "
            f"{dati[´parole´]} parole, "
            f"{dati[´caratteri´]} caratteri, "
            f"firma={dati[´firma´]}"
        )


if __name__ == "__main__":
    # Esempio: pipeline("dati_testo")
    pipeline("dati_testo")

Spiegazione

Il cuore dell’esempio è una pipeline a due stadi:

  • Stadio I/O: leggiamo più file contemporaneamente con ThreadPoolExecutor.
  • Stadio CPU: analizziamo i testi con ProcessPoolExecutor.

Perché usare i thread nella lettura?

La lettura dei file è un’operazione in cui il programma spesso resta in attesa del disco. In questi casi i thread sono utili perché mentre un thread aspetta I/O, un altro può continuare a lavorare. Non stiamo cercando di accelerare i calcoli, ma di nascondere la latenza.

Perché usare i processi nell’analisi?

La funzione analizza_testo() contiene un ciclo che elabora ogni carattere. Questo è un esempio di carico CPU-bound. In Python, il Global Interpreter Lock (GIL) limita l’esecuzione parallela dei thread per codice Python puro. Con i processi, invece, ogni worker ha il proprio interprete e può sfruttare più core della CPU.

Il flusso dei dati

  • La funzione trova_file_testo() individua i file da elaborare.
  • I thread leggono i contenuti in parallelo.
  • I risultati della lettura vengono raccolti in memoria come coppie (nome_file, testo).
  • I processi ricevono il testo e restituiscono un dizionario con statistiche e firma.
  • Infine i risultati vengono ordinati e stampati.

Dettagli importanti del codice

Ci sono alcuni punti tecnici da notare:

  • as_completed() permette di elaborare i task appena finiscono, senza aspettare l’ordine di invio.
  • max_workers nei thread è fissato a 4 come scelta ragionevole per I/O leggero; nei processi usiamo os.cpu_count() per sfruttare i core disponibili.
  • Nel blocco if __name__ == "__main__" evitiamo problemi di avvio dei processi, soprattutto su Windows e macOS.
  • La funzione passata ai processi deve essere serializzabile (picklable): per questo è definita a livello modulo.

Best practice

Questo tipo di architettura è molto efficace, ma va usata con criterio. Ecco le best practice più importanti.

  • Usa i thread per I/O, i processi per CPU: è la regola pratica più utile da ricordare.
  • Non creare troppi worker: aumentare indiscriminatamente thread o processi può peggiorare le prestazioni per overhead di scheduling e memoria.
  • Evita di passare oggetti enormi ai processi: ogni argomento deve essere serializzato. Se i dati sono molto grandi, valuta strategie diverse come chunking o file temporanei.
  • Isola bene le funzioni worker: funzioni pure, senza dipendenze globali complesse, sono più facili da testare e mantenere.
  • Gestisci sempre gli errori: un worker può fallire per file mancanti, encoding errato o dati malformati. Con future.result() le eccezioni emergono chiaramente e vanno intercettate.
  • Misura prima di ottimizzare: non tutto beneficia del multiprocessing. Su dataset piccoli, l’overhead può superare il guadagno.
  • Separa le fasi della pipeline: leggere, trasformare e analizzare sono responsabilità diverse. Questa separazione rende il codice più leggibile e scalabile.

Un consiglio pratico: se la tua applicazione fa sia I/O sia calcolo, spesso la soluzione migliore non è “solo thread” o “solo processi”, ma una combinazione ragionata dei due.

Riepilogo

In questo tutorial abbiamo visto un approccio molto concreto al tema multithreading e multiprocessing in Python: costruire una pipeline ibrida produttore-consumatore.

  • I thread sono ideali per leggere dati in parallelo quando il collo di bottiglia è l’I/O.
  • I processi sono la scelta giusta quando il lavoro è CPU-bound e vuoi superare il limite del GIL.
  • Una pipeline ben progettata ti permette di sfruttare entrambi i modelli in modo complementare.
  • Il codice resta più manutenibile se separi chiaramente le fasi di lettura e analisi.

Questo pattern è particolarmente utile in applicazioni reali come ETL, analisi di log, elaborazione di documenti, scraping con post-processing e pipeline di data processing.

Approfondisci con risorse ufficiali

SHARE