Introduzione ai DataFrame Dask
Sostituisca pd.read_csv e pd.DataFrame con gli equivalenti di dask, chiami compute() per avviare l'esecuzione e analizzi le prestazioni dei grafi di attività.
Introduzione ai DataFrame Dask è una lezione Pandas & NumPy Academy gratuita su CoddyKit. Questa è la lezione 3 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.
Che cos'è Dask
Dask è una libreria di calcolo parallelo per Python che estende NumPy e Pandas a dataset più grandi della RAM. Il modulo dask.dataframe fornisce un'API DataFrame quasi identica a quella di Pandas, ma invece di eseguire immediatamente le operazioni, Dask costruisce un grafo di attività e le esegue in modo differito quando si chiama .compute(). In questo modo Dask può parallelizzare il lavoro su più core o persino su più macchine, con modifiche minime al codice.
Installazione e importazione di Dask
Dask si installa con pip install dask[dataframe]. La convenzione per l'importazione è import dask.dataframe as dd. Internamente, un DataFrame Dask viene suddiviso in molti DataFrame Pandas più piccoli, ciascuno elaborato in modo indipendente. Le operazioni sul DataFrame Dask creano un grafo di attività differite: nulla viene eseguito finché non si chiama .compute(). Questa separazione tra la descrizione e l'esecuzione del calcolo è l'idea chiave di Dask.
import dask.dataframe as dd
# Read a large CSV — returns a Dask DataFrame immediately (no data loaded yet)
ddf = dd.read_csv('large_sales.csv')
print(type(ddf)) # dask.dataframe.core.DataFrame
print(ddf.columns.tolist())
print(ddf.dtypes)Dask e Pandas: la differenza fondamentale
Con Pandas, ogni operazione viene eseguita immediatamente e in modo eager. Con Dask, le operazioni restituiscono un altro oggetto Dask che rappresenta il calcolo differito. Solo quando si chiama .compute() Dask legge effettivamente i dati ed esegue il grafo di attività. Questa esecuzione differita consente a Dask di ottimizzare il piano prima dell'esecuzione: per esempio, può fondere filtri consecutivi per evitare di caricare i dati più volte. È come una ricetta: Dask scrive la ricetta, mentre .compute() prepara il piatto.
import dask.dataframe as dd
ddf = dd.read_csv('sales.csv')
# This does NOT run yet — just builds the task graph
filtered = ddf[ddf['amount'] > 1000]
agg = filtered.groupby('region')['amount'].sum()
print(type(agg)) # dask.dataframe.core.Series
# NOW execute everything
result = agg.compute()
print(result)Le partizioni: il concetto fondamentale
Un DataFrame Dask è suddiviso in partizioni, ognuna delle quali è un normale DataFrame Pandas. Per impostazione predefinita, dd.read_csv crea una partizione per file (oppure una partizione ogni 128 MB per i file di grandi dimensioni). È possibile controllare questo comportamento con blocksize. Verificando ddf.npartitions si visualizza il numero di partizioni esistenti. Un numero maggiore di partizioni aumenta il parallelismo, ma aggiunge overhead; un numero minore riduce l'overhead, ma limita il parallelismo. Il compromesso ideale consiste in genere in alcune centinaia di partizioni.
import dask.dataframe as dd
ddf = dd.read_csv('data/*.csv') # Read multiple CSV files at once
print('Number of partitions:', ddf.npartitions)
# Access a single partition as a Pandas DataFrame
first_partition = ddf.get_partition(0).compute()
print('Partition 0 shape:', first_partition.shape)Operazioni Pandas familiari in Dask
La maggior parte delle operazioni comuni di Pandas funziona allo stesso modo in Dask: .head(), .tail(), .describe(), l'indicizzazione booleana, .groupby(), .merge() e .assign() hanno tutte equivalenti in Dask. La differenza principale è che è necessario chiamare .compute() per materializzare il risultato. Le operazioni che Pandas gestisce in millisecondi possono richiedere secondi in Dask a causa dell'overhead del grafo di attività: utilizzi quindi Pandas per i dati di piccole dimensioni e Dask quando i dati non entrano nella RAM.
import dask.dataframe as dd
ddf = dd.read_csv('transactions.csv')
# Filtering — same syntax as Pandas
high_value = ddf[ddf['amount'] > 500]
# GroupBy aggregation
by_region = high_value.groupby('region')['amount'].mean()
# Execute
result = by_region.compute()
print(result.sort_values(ascending=False))Lettura di più file con i pattern glob
Una delle funzionalità più utili di Dask è la lettura di più file contemporaneamente tramite pattern glob. dd.read_csv('data/2024-*.csv') legge tutti i file corrispondenti e crea una partizione per file. È ideale per i dati archiviati in file suddivisi per mese o per giorno, uno schema comune nei data lake. Dask allinea automaticamente gli schemi: è l'equivalente di un ciclo manuale con concatenazione in Pandas, ma molto più semplice.
import dask.dataframe as dd
# Read all monthly files at once
ddf = dd.read_csv('sales/2024-*.csv',
dtype={'order_id': 'int32',
'amount': 'float32'})
print(f'Partitions: {ddf.npartitions}') # One per file
print(f'Total rows (lazy): {len(ddf)}') # This triggers a compute!Il metodo visualize() per i grafi di attività
Prima di eseguire una pipeline Dask complessa, è possibile esaminare il grafo di attività chiamando result.visualize(), che genera un diagramma PNG di tutti i passaggi del calcolo. È utile per capire cosa eseguirà Dask e per individuare le cause di rallentamenti imprevisti. Il grafo mostra come le partizioni passano attraverso i passaggi di filtraggio, raggruppamento e aggregazione, rendendo facile individuare i calcoli ridondanti. È richiesto il pacchetto graphviz.
import dask.dataframe as dd
ddf = dd.read_csv('orders.csv')
pipeline = (
ddf[ddf['status'] == 'completed']
.groupby('product_id')['revenue']
.sum()
)
# Visualise the task graph (saves to PNG)
# pipeline.visualize('task_graph.png')
# Check number of tasks in the graph
print('Number of tasks:', len(pipeline.__dask_graph__()))Applicazione di funzioni personalizzate con map_partitions
Quando è necessario applicare una funzione Pandas personalizzata a un DataFrame Dask, si utilizzi ddf.map_partitions(func). Questo applica func a ogni partizione in modo indipendente e restituisce un nuovo DataFrame Dask. La funzione riceve un normale DataFrame Pandas e deve restituirne uno. È l'equivalente Dask di df.apply() e consente di integrare Dask con codice che comprende soltanto Pandas.
import dask.dataframe as dd
import pandas as pd
def normalise_chunk(df):
df = df.copy()
df['amount_norm'] = (df['amount'] - df['amount'].mean()) / df['amount'].std()
return df
ddf = dd.read_csv('data.csv')
normalised = ddf.map_partitions(normalise_chunk)
result = normalised[['order_id', 'amount_norm']].compute()
print(result.head())Opzioni dello scheduler di Dask
Dask dispone di più scheduler che controllano il modo in cui vengono eseguite le attività. Lo scheduler 'synchronous' esegue le attività in sequenza nel thread corrente (utile per il debugging). Lo scheduler 'threads' utilizza un pool di thread (adatto alle operazioni vincolate dall'I/O). Lo scheduler 'processes' avvia più processi per le operazioni vincolate dalla CPU (aggirando il GIL di Python). Un cluster Dask distribuito consente l'esecuzione su più macchine. Specifichi lo scheduler tramite compute(scheduler='threads').
import dask.dataframe as dd
ddf = dd.read_csv('data.csv')
agg = ddf.groupby('category')['sales'].sum()
# Choose scheduler based on workload
result_sync = agg.compute(scheduler='synchronous') # sequential, easy to debug
result_threads = agg.compute(scheduler='threads') # parallel I/O
result_processes = agg.compute(scheduler='processes') # parallel CPUConversione tra Dask e Pandas
È comune elaborare un dataset di grandi dimensioni con Dask e poi trasferire il risultato aggregato in Pandas per l'analisi o la visualizzazione finale. Utilizzi .compute() per convertire un DataFrame Dask in Pandas. Nella direzione opposta, dd.from_pandas(df, npartitions=4) converte un DataFrame Pandas in un DataFrame Dask, utile per testare il codice Dask su dati di piccole dimensioni prima di estenderlo all'intero dataset.
import pandas as pd
import dask.dataframe as dd
# Start with a small Pandas DF for testing
df_small = pd.DataFrame({'a': range(100), 'b': range(100, 200)})
# Convert to Dask for development/testing
ddf = dd.from_pandas(df_small, npartitions=4)
result = ddf.groupby('a')['b'].sum().compute()
print(type(result)) # pandas.Series
print(result.head())Quando usare Dask, Pandas o SQL
Dask non è sempre lo strumento più adatto. Utilizzi Pandas quando i dati entrano nella RAM (meno di qualche GB): è più semplice e veloce grazie al minore overhead. Utilizzi Dask quando i dati superano la RAM, ma desidera una sintassi simile a Pandas e il parallelismo su una singola macchina. Utilizzi SQL o un database quando i dati risiedono in un database relazionale e le aggregazioni possono essere delegate al motore del database. Utilizzi Spark o BigQuery quando è necessaria un'elaborazione distribuita su più macchine, su scala petabyte.
Verifica rapida
Verifichi la propria comprensione dei concetti di analisi dei dati trattati in questa lezione.
Riepilogo della lezione
In questa lezione ha imparato che i DataFrame Dask sono raccolte di partizioni Pandas valutate in modo differito e dotate di un'API familiare, che .compute() avvia l'esecuzione effettiva del grafo di attività e che map_partitions consente di applicare qualsiasi funzione Pandas personalizzata a tutte le partizioni. Prossimamente esamineremo il formato Parquet, un'alternativa colonnare e veloce al CSV per l'archiviazione di dataset di grandi dimensioni.
Domande Frequenti
La lezione «Introduzione ai DataFrame Dask» è gratuita?
Sì — il testo completo di «Introduzione ai DataFrame Dask» è 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 «Introduzione ai DataFrame Dask»?
Sostituisca pd.read_csv e pd.DataFrame con gli equivalenti di dask, chiami compute() per avviare l'esecuzione e analizzi le prestazioni dei grafi di attività. 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 3 di 4.
Quanto tempo richiede la lezione «Introduzione ai DataFrame Dask»?
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
- CSV in streaming con chunksize
- Aggregazione incrementale tra i blocchi
- Introduzione ai DataFrame Dask
- Parquet: archiviazione colonnare veloce