Pandas & NumPy Academy · Lezione

Pianificare e registrare le esecuzioni della pipeline

Esegua la pipeline come script Python dalla riga di comando, registri gli orari di inizio e fine e usi cron o uno scheduler per automatizzarla.

Lezione 4 di 413 passaggi

Pianificare e registrare le esecuzioni della pipeline è una lezione Pandas & NumPy Academy gratuita su CoddyKit. Questa è la lezione 4 di 4. Puoi leggere la lezione completa qui gratuitamente — poi esercitati direttamente nel browser con un editor di codice integrato e un tutor IA disponibile 24/7. Fa parte del percorso di apprendimento Pandas & NumPy Academy, e i tuoi progressi si sincronizzano tra il web e l'app CoddyKit. Il corso Pandas & NumPy Academy include 4 lezioni in totale.

Dal notebook allo script

Una pipeline che viene eseguita solo quando uno sviluppatore apre manualmente un notebook non offre alcun valore aziendale oltre la prima esecuzione. Per eseguirla automaticamente ogni giorno, la pipeline deve essere strutturata come uno script Python eseguibile dalla riga di comando: python pipeline.py. Questo richiede un punto di ingresso if __name__ == '__main__':, l'analisi degli argomenti della riga di comando e un sistema di logging adeguato: i tre pilastri di uno script pronto per la produzione.

# pipeline.py
import argparse
import logging
import pandas as pd

def main(config_path):
    logging.info(f'Starting pipeline with config: {config_path}')
    # ... run ETL steps ...
    logging.info('Pipeline complete.')

if __name__ == '__main__':
    parser = argparse.ArgumentParser()
    parser.add_argument('--config', default='config.json')
    args = parser.parse_args()
    main(args.config)

Configurazione del logging in Python

Il modulo integrato logging di Python è lo strumento corretto per i log della pipeline, non le istruzioni print(). Configuri un logger con output sia sulla console sia su file usando logging.basicConfig(). Usi il livello INFO per l'avanzamento normale e ERROR per gli errori. I log salvati su file persistono dopo la terminazione del processo, una caratteristica essenziale per il debug delle esecuzioni pianificate che nessuno sta osservando.

import logging
from datetime import date

log_file = f'pipeline_{date.today()}.log'

logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s %(levelname)s %(message)s',
    handlers=[
        logging.FileHandler(log_file),
        logging.StreamHandler()
    ]
)

logging.info('Logger configured.')

Registrazione dell'inizio e della fine della pipeline

Registri sempre l'ora di inizio, l'ora di fine e il tempo trascorso di un'esecuzione della pipeline. In questo modo stabilisce una baseline: se normalmente la pipeline viene eseguita in 45 secondi e oggi ha impiegato 8 minuti, qualcosa è cambiato, ad esempio il file di input è 10 volte più grande oppure una query al database è lenta. Le voci di log con timestamp relative all'inizio e alla fine rendono il confronto immediato basandosi solo sul file di log.

import time
import logging

def run_pipeline(config):
    start = time.time()
    logging.info(f'Pipeline START | env={config.get("env", "dev")} | input={config["input_path"]}')

    try:
        df = extract(config)
        df_clean = transform(df, config)
        load(df_clean, config)
        elapsed = time.time() - start
        logging.info(f'Pipeline SUCCESS | rows={len(df_clean)} | elapsed={elapsed:.1f}s')
    except Exception as e:
        logging.error(f'Pipeline FAILED | error={e}', exc_info=True)
        raise

Registrazione del numero di righe per ogni passaggio

Registri il numero di righe in ingresso e in uscita per ogni passaggio di trasformazione. Un log chiaro è simile a questo: extract: 50,000 rows → drop_nulls: 49,200 rows → filter: 47,800 rows → output: 47,800 rows. Questa traccia rende immediatamente evidente quante righe sono state eliminate in ogni passaggio e se i numeri sono quelli attesi. Riduzioni anomale emergono come variazioni nei conteggi registrati.

def log_step(df, step_name):
    logging.info(f'{step_name}: {len(df):,} rows')
    return df

import pandas as pd

df = (pd.read_csv('orders.csv')
    .pipe(log_step, 'extract')
    .dropna(subset=['revenue'])
    .pipe(log_step, 'drop_nulls')
    .query('quantity > 0')
    .pipe(log_step, 'filter_qty')
)
print('Step logging complete.')

Pianificazione con cron su Linux/Mac

cron è lo scheduler Unix standard per i processi ricorrenti. Modifichi il crontab con crontab -e e aggiunga una riga che specifichi quando eseguire lo script. Il formato è: minuto ora giorno mese giorno-della-settimana comando. Per eseguire una pipeline ogni giorno alle 06:00, usi 0 6 * * * /usr/bin/python /path/to/pipeline.py. Usi sempre percorsi assoluti nelle voci cron, perché cron viene eseguito in un ambiente minimo, privo delle impostazioni PATH della shell.

# crontab entry — edit with: crontab -e
# Run pipeline.py at 06:00 every day
# 0 6 * * * /opt/homebrew/bin/python /Users/analyst/pipeline.py --config /Users/analyst/config.json >> /Users/analyst/cron.log 2>&1

# Common cron patterns:
# 0 6 * * *      — daily at 06:00
# 0 */4 * * *    — every 4 hours
# 0 9 * * 1      — every Monday at 09:00
print('Cron schedule format: minute hour day month weekday')

Pianificazione con la libreria Python schedule

La libreria schedule offre un modo interamente basato su Python per eseguire processi a intervalli specifici senza usare cron. È utile negli ambienti in cui cron non è disponibile (Windows) o quando desidera che la logica dello scheduler risieda direttamente nel processo Python. Inserisca la pipeline in un ciclo di processi pianificati e mantenga il processo attivo per eseguire ripetutamente il lavoro.

# pip install schedule
# import schedule, time

# def job():
#     logging.info('Scheduled run starting...')
#     run_pipeline(CONFIG)

# schedule.every().day.at('06:00').do(job)
# schedule.every(4).hours.do(job)

# while True:
#     schedule.run_pending()
#     time.sleep(60)

print('schedule library: use for in-process Python scheduling')

Gestione degli errori e codici di uscita

Quando si verifica un errore, uno script di pipeline dovrebbe restituire un codice di uscita diverso da zero, così lo scheduler sa che il processo non è riuscito. Inserisca l'esecuzione principale in un blocco try/except e chiami sys.exit(1) in caso di errore. cron, Jenkins e Airflow controllano tutti il codice di uscita: un codice diverso da zero attiva un avviso, una nuova esecuzione o una notifica. Un'eccezione non gestita che non imposta il codice di uscita potrebbe passare inosservata ai sistemi di monitoraggio automatizzati.

import sys

def main():
    try:
        run_pipeline(CONFIG)
        sys.exit(0)  # success
    except AssertionError as e:
        logging.error(f'Data validation failed: {e}')
        sys.exit(2)  # data error
    except Exception as e:
        logging.error(f'Unexpected error: {e}', exc_info=True)
        sys.exit(1)  # general failure

print('Exit code 0=success, 1=error, 2=data failure')

Scrittura di un file di riepilogo dell'esecuzione della pipeline

Dopo un'esecuzione riuscita, scriva un piccolo file di riepilogo JSON insieme all'output. Includa il timestamp dell'esecuzione, il numero di righe in ingresso, il numero di righe in uscita, le righe eliminate e il tempo trascorso. I sistemi di monitoraggio e i dashboard possono leggere questo file per seguire nel tempo le tendenze relative allo stato della pipeline. Un dashboard che mostra le righe di output negli ultimi 30 giorni rende facile individuare il giorno in cui una sorgente dati ha iniziato a fornire meno record.

import json
from datetime import datetime

def write_run_summary(config, input_rows, output_rows, elapsed):
    summary = {
        'run_at': datetime.now().isoformat(),
        'input_path': config['input_path'],
        'input_rows': input_rows,
        'output_rows': output_rows,
        'rows_dropped': input_rows - output_rows,
        'elapsed_seconds': round(elapsed, 2),
        'status': 'success'
    }
    with open('last_run_summary.json', 'w') as f:
        json.dump(summary, f, indent=2)
    print('Run summary written.')

Pianificazione idempotente: evitare le doppie esecuzioni

Se una pipeline pianificata viene attivata accidentalmente due volte, non dovrebbe danneggiare l'output. Progetti il passaggio di caricamento in modo idempotente: usi un nome file di output con la data oppure sovrascriva lo stesso output con il risultato più recente. Per i caricamenti nei database, usi if_exists='replace' oppure un approccio UPSERT. Non usi mai la modalità append senza un passaggio di deduplicazione, altrimenti ogni esecuzione pianificata aggiungerà righe duplicate alla tabella di output.

from datetime import date

def load_idempotent(df, config):
    # Date-stamped output: each run overwrites its own day's file
    output_path = f"output_{date.today().strftime('%Y%m%d')}.parquet"
    df.to_parquet(output_path, index=False)
    logging.info(f'Loaded {len(df)} rows to {output_path}')

Avvisi in caso di errore della pipeline

Per le pipeline da cui dipendono le operazioni aziendali, il silenzio dopo un errore è pericoloso. Configuri un avviso semplice: se il file di riepilogo dell'esecuzione non viene aggiornato entro l'intervallo previsto, invii un'e-mail o un messaggio Slack. smtplib di Python può inviare un'e-mail in caso di errore, oppure può usare un webhook per pubblicare un messaggio su Slack. Generi immediatamente un avviso in caso di codice di uscita 1 o 2, così l'analista sa che l'aggiornamento giornaliero non è riuscito prima che se ne accorga l'azienda.

import smtplib

def send_failure_alert(error_msg):
    # Example: send plain-text email via SMTP
    # server = smtplib.SMTP('smtp.example.com', 587)
    # server.sendmail('pipeline@company.com',
    #                 'analyst@company.com',
    #                 f'Subject: Pipeline Failed\n\n{error_msg}')
    # server.quit()
    print(f'[ALERT] Would send failure notification: {error_msg}')

# In main():
# except Exception as e:
#     send_failure_alert(str(e))
#     sys.exit(1)
print('Alert integration pattern shown above.')

Script completo per una pipeline pianificata

Unisca tutti gli elementi — analisi degli argomenti, configurazione del logging, riepilogo dell'esecuzione, gestione degli errori e codici di uscita — in uno script completo per la pipeline. Questo script può essere inserito in qualsiasi ambiente, collegato a un file di configurazione e pianificato con cron o con qualsiasi orchestratore di workflow. A ogni esecuzione produce un file di log datato, un riepilogo dell'esecuzione e un file di output datato, rendendo ogni esecuzione completamente verificabile e riproducibile in modo indipendente.

# Full script skeleton:
# 1. parse --config argument
# 2. configure logging to file + console
# 3. load JSON config
# 4. validate config
# 5. run extract() -> transform() -> load()
# 6. write run summary JSON
# 7. sys.exit(0) on success, sys.exit(1) on failure

print('Production pipeline script structure complete.')
print('Schedule with: crontab -e  or  python scheduler.py')

Verifica rapida

Verifichi la Sua comprensione dei concetti di analisi dei dati presentati in questa lezione.

Riepilogo della lezione

In questa lezione ha imparato a: strutturare una pipeline come script da riga di comando con analisi degli argomenti e logging, pianificare l'esecuzione con cron e gestire gli errori tramite codici di uscita diversi da zero e avvisi, nonché scrivere file di riepilogo delle esecuzioni e progettare fasi di caricamento idempotenti per un'esecuzione automatizzata affidabile. Congratulazioni per aver completato il percorso Analisi dei dati: Pandas e NumPy!

Gratis per iniziare

Impara Python con un tutor IA — gratis

Scrivi ed esegui vero codice nel tuo browser, ricevi aiuto istantaneo da un tutor IA disponibile 24/7, e riprendi da dove hai lasciato sul web o nell'app.

Corsi
30
Lezioni
120

Domande Frequenti

La lezione «Pianificare e registrare le esecuzioni della pipeline» è gratuita?

Sì — il testo completo di «Pianificare e registrare le esecuzioni della pipeline» è gratuito qui sul web. Per esercitarvi in modo interattivo (un editor di codice integrato e un tutor IA 24/7) e sbloccare il resto del corso Pandas & NumPy Academy, passa a CoddyKit PRO. Il corso Pandas & NumPy Academy include 4 lezioni in totale.

Cosa imparerò in «Pianificare e registrare le esecuzioni della pipeline»?

Esegua la pipeline come script Python dalla riga di comando, registri gli orari di inizio e fine e usi cron o uno scheduler per automatizzarla. Eserciti Pandas & NumPy Academy con codice pratico che esegui direttamente nel browser, e un tutor IA 24/7 risponde alle tue domande mentre lavori sulla lezione.

Ho bisogno di esperienza per iniziare Pandas & NumPy Academy?

Non è richiesta alcuna esperienza precedente. Pandas & NumPy Academy su CoddyKit è strutturato per principianti e studenti avanzati, quindi puoi iniziare da qui o dall'inizio e procedere al tuo ritmo. Questa è la lezione 4 di 4.

Quanto tempo richiede la lezione «Pianificare e registrare le esecuzioni della pipeline»?

La maggior parte delle lezioni CoddyKit richiede circa 5–10 minuti. Ogni lezione è breve e interattiva, quindi fai progressi costanti e riprendi esattamente da dove hai lasciato su web e app.

Posso scrivere ed eseguire codice in questa lezione Pandas & NumPy Academy?

Sì. Ogni lezione Pandas & NumPy Academy include un editor di codice integrato, quindi scrivi ed esegui codice reale direttamente nel tuo browser e ricevi feedback istantaneo dall'IA — nessuna configurazione locale necessaria.

Tutte le lezioni di questo corso

  1. Strutturare i passaggi di trasformazione come funzioni
  2. Parametrizzare le pipeline con dict di configurazione
  3. Testare i passaggi della pipeline con asserzioni
  4. Pianificare e registrare le esecuzioni della pipeline
← Torna a Pandas & NumPy Academy