Aggregazione incrementale tra i blocchi
Accumuli conteggi, somme e valori minimi/massimi progressivi tra i blocchi senza conservare in memoria l'intero file.
Aggregazione incrementale tra i blocchi è una lezione Pandas & NumPy Academy gratuita su CoddyKit. Questa è la lezione 2 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.
Perché usare l'aggregazione incrementale?
L'aggregazione incrementale è fondamentale per analizzare dataset più grandi della RAM senza distribuire i calcoli su più macchine. Invece di caricare tutti i dati per calcolare una statistica finale, si mantengono accumulatori progressivi — somme parziali, conteggi, valori minimi/massimi — aggiornandoli a ogni blocco. Il risultato finale viene assemblato a partire da questi accumulatori leggeri dopo aver scansionato l'intero file. Questo modello permette di gestire file di terabyte su un singolo laptop.
Contare le righe e calcolare la media
Per calcolare la media tra più blocchi è necessario tenere traccia separatamente della somma e del conteggio progressivi. Non è possibile calcolare semplicemente la media delle medie dei blocchi, perché i blocchi possono avere dimensioni diverse. La formula corretta è total_sum / total_count. Questo modello si estende a qualsiasi quantità scomponibile: varianza, correlazione e istogrammi dispongono tutti di formule incrementali.
import pandas as pd
total_sum = 0.0
total_count = 0
for chunk in pd.read_csv('transactions.csv', chunksize=100000):
total_sum += chunk['amount'].sum()
total_count += chunk['amount'].notna().sum()
grand_mean = total_sum / total_count
print(f'Rows processed: {total_count:,}')
print(f'Grand mean: {grand_mean:.4f}')Minimo e massimo incrementali
Tenere traccia del minimo e del massimo globali tra i blocchi è semplice: inizializzi i valori con float('inf') e float('-inf') di Python, quindi li aggiorni con il minimo/massimo di ogni blocco. In questo modo si evita qualsiasi memorizzazione intermedia. Il modello si generalizza al minimo/massimo per gruppo mantenendo un dizionario indicizzato dall'identificatore del gruppo.
import pandas as pd
global_min = float('inf')
global_max = float('-inf')
for chunk in pd.read_csv('prices.csv', chunksize=50000):
chunk_min = chunk['price'].min()
chunk_max = chunk['price'].max()
if chunk_min < global_min:
global_min = chunk_min
if chunk_max > global_max:
global_max = chunk_max
print(f'Price range: {global_min} to {global_max}')Conteggi incrementali delle frequenze
Per le colonne categoriali, mantenga un dizionario progressivo delle frequenze aggiungendo il risultato di value_counts() di ogni blocco a un accumulatore Pandas Series. Poiché l'addizione delle Pandas Series allinea le etichette degli indici, le categorie sconosciute presenti nei blocchi successivi vengono incluse automaticamente. Dopo aver elaborato tutti i blocchi, ordini i risultati per conteggio per visualizzare le categorie più frequenti nell'intero dataset.
import pandas as pd
freq = pd.Series(dtype='int64')
for chunk in pd.read_csv('orders.csv',
chunksize=100000,
usecols=['category']):
chunk_counts = chunk['category'].value_counts()
freq = freq.add(chunk_counts, fill_value=0)
# Final sorted frequency table
print(freq.sort_values(ascending=False).head(10))Aggregazione GroupBy incrementale
Per calcolare una somma o un conteggio groupby tra più blocchi, applichi groupby().agg() all'interno di ogni blocco e memorizzi la Series o il DataFrame risultante. Dopo il ciclo, concateni tutti i risultati parziali e applichi un secondo groupby per combinarli. Questo approccio in due fasi gestisce correttamente i gruppi presenti in più blocchi, una situazione comune quando i dati sono ordinati per data anziché per gruppo.
import pandas as pd
partials = []
for chunk in pd.read_csv('sales.csv',
chunksize=100000,
usecols=['region', 'product', 'revenue']):
p = chunk.groupby(['region', 'product'])['revenue'].sum()
partials.append(p)
final = (
pd.concat(partials)
.groupby(level=['region', 'product'])
.sum()
.sort_values(ascending=False)
)
print(final.head(10))Calcolare la varianza in modo incrementale (metodo di Welford)
Calcolare la varianza tra più blocchi è più complesso che calcolare la media. La formula ingenua E[X²] - E[X]² soffre di cancellazione catastrofica in presenza di medie elevate. L'algoritmo online di Welford mantiene una media progressiva e la somma degli scarti quadratici, aggiornandoli con ogni nuovo valore in modo numericamente stabile. Sebbene SciPy lo implementi, comprendere il modello consente di estenderlo alla varianza ponderata e alla covarianza.
import pandas as pd
import numpy as np
# Simple two-pass approach using stored chunk stats
chunk_stats = []
for chunk in pd.read_csv('data.csv',
chunksize=100000,
usecols=['value']):
n = chunk['value'].count()
mean = chunk['value'].mean()
var = chunk['value'].var(ddof=1)
chunk_stats.append((n, mean, var))
# Combine: use pooled variance formula
total_n = sum(s[0] for s in chunk_stats)
total_mean = sum(s[0]*s[1] for s in chunk_stats) / total_n
pooled_var = sum((s[0]-1)*s[2] + s[0]*(s[1]-total_mean)**2
for s in chunk_stats) / (total_n - 1)
print(f'Grand variance: {pooled_var:.4f}')Creare un istogramma incrementale
Calcolare la distribuzione di una colonna in un file troppo grande per la RAM richiede un istogramma incrementale. Stabilizzi in anticipo i limiti degli intervalli (in base a un piccolo campione o alla conoscenza del dominio), quindi utilizzi np.histogram(chunk_values, bins=edges) all'interno di ogni blocco e accumuli i conteggi. Al termine, rappresenti i conteggi combinati in un grafico a barre. È così che sistemi di streaming come Kafka Streams e Flink calcolano istogrammi approssimati.
import pandas as pd
import numpy as np
# Decide bin edges from a sample
sample = pd.read_csv('amounts.csv', nrows=5000)
bins = np.linspace(sample['amount'].min(),
sample['amount'].max(), 21)
counts = np.zeros(len(bins) - 1, dtype='int64')
for chunk in pd.read_csv('amounts.csv',
chunksize=100000,
usecols=['amount']):
chunk_counts, _ = np.histogram(
chunk['amount'].dropna(), bins=bins
)
counts += chunk_counts
print('Histogram counts:', counts[:5], '...')Tracciare in modo approssimato i valori univoci
Contare i valori distinti esatti tra più blocchi richiede di memorizzare tutti i valori univoci, potenzialmente milioni. Per conteggi approssimati su larga scala, utilizzi uno sketch HyperLogLog, disponibile in Python tramite la libreria hyperloglog. In alternativa, tracci gli elementi univoci di ogni blocco con un set e ne calcoli l'unione, ma questa struttura cresce senza limiti. Per un'approssimazione economica, utilizzi pd.Series.nunique() per ogni blocco e riporti la media: non è esatta, ma spesso è sufficiente per la profilazione dei dati.
import pandas as pd
unique_ids = set()
for chunk in pd.read_csv('events.csv',
chunksize=100000,
usecols=['user_id']):
unique_ids.update(chunk['user_id'].dropna().unique())
print(f'Distinct user IDs: {len(unique_ids):,}')
# Warning: the set may grow large for high-cardinality columnsSegnalare l'avanzamento durante le esecuzioni prolungate
L'elaborazione di un file di diversi gigabyte può richiedere alcuni minuti. Aggiunga una segnalazione dell'avanzamento per sapere che la pipeline è in esecuzione e poter stimare il tempo rimanente. Conti i byte o le righe elaborati e li confronti con le dimensioni del file. La libreria tqdm rende questa operazione molto semplice grazie al wrapper tqdm(reader). Anche senza tqdm, stampare una riga di stato ogni 10 blocchi fornisce un riscontro utile durante i processi batch prolungati.
import pandas as pd
import time
chunksize = 100000
start = time.time()
rows_processed = 0
for i, chunk in enumerate(pd.read_csv('big.csv',
chunksize=chunksize)):
rows_processed += len(chunk)
# Report every 10 chunks
if (i + 1) % 10 == 0:
elapsed = time.time() - start
rate = rows_processed / elapsed
print(f'Chunk {i+1}: {rows_processed:,} rows '
f'@ {rate/1000:.0f}k rows/sec')
print(f'Total: {rows_processed:,} rows in {time.time()-start:.1f}s')Filtrare prima di aggregare
Applichi i filtri all'interno di ogni blocco prima dell'aggregazione, per evitare di accumulare dati indesiderati. Ad esempio, se Le interessano solo gli ordini del 2024, filtri la colonna della data del blocco prima di eseguire il groupby. In questo modo riduce la memoria necessaria per i risultati parziali e velocizza il passaggio finale di concatenazione. Inserisca sempre i filtri il più presto possibile nella pipeline: è un principio fondamentale dell'elaborazione efficiente dei dati.
import pandas as pd
partials = []
for chunk in pd.read_csv('orders.csv',
chunksize=100000,
parse_dates=['order_date']):
# Filter early: only 2024 orders
mask = chunk['order_date'].dt.year == 2024
filtered = chunk.loc[mask, ['category', 'revenue']]
if len(filtered) > 0:
p = filtered.groupby('category')['revenue'].sum()
partials.append(p)
if partials:
result = pd.concat(partials).groupby(level=0).sum()
print(result)Salvare i risultati intermedi
Per i processi di durata molto lunga, salvi periodicamente i risultati intermedi in modo da poter riprendere da un checkpoint se il processo viene interrotto. Scriva gli aggregati per blocco in un file Parquet o CSV dopo ogni N blocchi. Se il processo si interrompe al blocco 800 su 1000, può ricaricare gli aggregati salvati e continuare dal punto raggiunto, invece di rielaborare l'intero file. Questo schema di resilienza è essenziale nelle pipeline di dati per la produzione.
import pandas as pd
import os
CHECKPOINT = 'checkpoint.csv'
running_total = 0.0
running_count = 0
# Resume from checkpoint if it exists
if os.path.exists(CHECKPOINT):
ckpt = pd.read_csv(CHECKPOINT)
running_total = ckpt['total'].iloc[0]
running_count = int(ckpt['count'].iloc[0])
print(f'Resuming from checkpoint: {running_count:,} rows')
for chunk in pd.read_csv('huge.csv', chunksize=100000):
running_total += chunk['value'].sum()
running_count += len(chunk)
# Save checkpoint
pd.DataFrame({'total': [running_total],
'count': [running_count]}).to_csv(CHECKPOINT, index=False)
print(f'Final mean: {running_total / running_count:.4f}')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: gli accumulatori incrementali (somma, conteggio, minimo/massimo, Series delle frequenze) consentono di eseguire aggregazioni con memoria costante su file di grandi dimensioni; il groupby in due fasi (groupby parziale per blocco, seguito da concatenazione e nuovo raggruppamento) gestisce correttamente i gruppi distribuiti tra più blocchi; e il filtraggio anticipato all'interno di ogni blocco riduce il costo della fase di accumulo. Prossimamente esploreremo i DataFrame Dask come sostituti paralleli immediati di Pandas per i dataset di grandi dimensioni.
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 «Aggregazione incrementale tra i blocchi» è gratuita?
Sì — il testo completo di «Aggregazione incrementale tra i blocchi» è 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 «Aggregazione incrementale tra i blocchi»?
Accumuli conteggi, somme e valori minimi/massimi progressivi tra i blocchi senza conservare in memoria l'intero file. 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 2 di 4.
Quanto tempo richiede la lezione «Aggregazione incrementale tra i blocchi»?
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