Inkrementelle Aggregation über Blöcke
Akkumulieren Sie laufende Anzahlen, Summen sowie Minima und Maxima über mehrere Blöcke, ohne die vollständige Datei im Speicher abzulegen.
Inkrementelle Aggregation über Blöcke ist eine kostenlose Pandas & NumPy Academy-Lektion auf CoddyKit. Dies ist Lektion 2 von 4. Du kannst die komplette Lektion unten kostenlos lesen – dann übst du sie direkt im Browser mit einem integrierten Code-Editor und einem KI-Tutor rund um die Uhr. Sie ist Teil des Pandas & NumPy Academy-Lernpfads, und dein Fortschritt wird über Web und CoddyKit-App synchronisiert. Der Pandas & NumPy Academy-Kurs umfasst insgesamt 4 Lektionen.
Warum inkrementell aggregieren?
Inkrementelle Aggregation ist der Schlüssel zur Analyse von Datensätzen, die größer als der Arbeitsspeicher sind, ohne die Berechnung auf mehrere Rechner zu verteilen. Statt alle Daten zu laden, um eine endgültige Statistik zu berechnen, verwalten Sie laufende Akkumulatoren — Teilsummen, Anzahlen sowie Minimal- und Maximalwerte — und aktualisieren sie mit jedem Block. Das Endergebnis wird aus diesen kompakten Akkumulatoren zusammengestellt, nachdem die gesamte Datei eingelesen wurde. Mit diesem Muster lassen sich Terabyte große Dateien auf einem einzigen Laptop verarbeiten.
Zeilen zählen und den Mittelwert berechnen
Für die Berechnung des Mittelwerts über mehrere Blöcke müssen die laufende Summe und die Anzahl getrennt verfolgt werden. Sie können nicht einfach die Mittelwerte der einzelnen Blöcke mitteln, da die Blöcke unterschiedlich groß sein können. Die korrekte Formel lautet total_sum / total_count. Dieses Muster lässt sich auf jede Größe übertragen, die sich zerlegen lässt: Für Varianz, Korrelation und Histogramme gibt es ebenfalls inkrementelle Formeln.
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}')Inkrementelles Minimum und Maximum
Das globale Minimum und Maximum über mehrere Blöcke hinweg zu verfolgen, ist unkompliziert: Initialisieren Sie die Werte mit Pythons float('inf') und float('-inf') und aktualisieren Sie sie anschließend mit dem Minimum bzw. Maximum jedes Blocks. Dadurch ist keine Zwischenspeicherung erforderlich. Das Muster lässt sich auf Minimum und Maximum pro Gruppe verallgemeinern, indem ein nach Gruppenbezeichnern geordnetes Wörterbuch verwaltet wird.
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}')Inkrementelle Häufigkeitszählungen
Für kategoriale Spalten verwalten Sie ein laufendes Häufigkeitswörterbuch, indem Sie das Ergebnis von value_counts() jedes Blocks zu einem Pandas-Series-Akkumulator addieren. Da die Addition von Pandas Series anhand der Indexbezeichnungen ausgerichtet wird, werden unbekannte Kategorien in späteren Blöcken automatisch aufgenommen. Sortieren Sie nach der Verarbeitung aller Blöcke nach der Anzahl, um die häufigsten Kategorien im gesamten Datensatz zu sehen.
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))Inkrementelle GroupBy-Aggregation
Um eine groupby-Summe oder -Anzahl über mehrere Blöcke hinweg zu berechnen, wenden Sie groupby().agg() innerhalb jedes Blocks an und speichern die resultierende Series oder den resultierenden DataFrame. Verketten Sie nach der Schleife alle Teilergebnisse und wenden Sie eine zweite groupby-Aggregation an, um sie zusammenzuführen. Dieser zweistufige Ansatz behandelt Gruppen, die in mehreren Blöcken vorkommen, korrekt. Das ist häufig der Fall, wenn die Daten nach Datum statt nach Gruppe sortiert sind.
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))Varianz inkrementell berechnen (Welford-Methode)
Die Berechnung der Varianz über mehrere Blöcke hinweg ist schwieriger als die des Mittelwerts. Die naive Formel E[X²] - E[X]² leidet bei großen Mittelwerten unter katastrophaler Auslöschung. Welfords Online-Algorithmus verwaltet einen laufenden Mittelwert und die Summe der quadrierten Abweichungen und aktualisiert beide Werte mit jedem neuen Datenpunkt auf numerisch stabile Weise. SciPy implementiert dieses Verfahren zwar, aber das Verständnis des Musters ermöglicht es Ihnen, es auf gewichtete Varianz und Kovarianz zu erweitern.
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}')Ein inkrementelles Histogramm erstellen
Die Verteilung einer Spalte über eine Datei zu berechnen, die zu groß für den Arbeitsspeicher ist, erfordert ein inkrementelles Histogramm. Legen Sie die Klassengrenzen im Voraus fest (auf Grundlage einer kleinen Stichprobe oder von Domänenwissen), verwenden Sie dann innerhalb jedes Blocks np.histogram(chunk_values, bins=edges) und akkumulieren Sie die Anzahlen. Zeichnen Sie am Ende die kombinierten Anzahlen als Balkendiagramm. So berechnen Streaming-Systeme wie Kafka Streams und Flink Näherungshistogramme.
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], '...')Eindeutige Werte näherungsweise verfolgen
Für die exakte Anzahl unterschiedlicher Werte über mehrere Blöcke hinweg müssten alle eindeutigen Werte gespeichert werden — potenziell mehrere Millionen. Für Näherungswerte im großen Maßstab verwenden Sie einen HyperLogLog-Sketch, der in Python über die Bibliothek hyperloglog verfügbar ist. Alternativ können Sie die eindeutigen Werte pro Block in einer Menge erfassen und die Vereinigungsmenge bilden, diese wächst jedoch unbegrenzt. Für eine kostengünstige Näherung können Sie pd.Series.nunique() pro Block verwenden und den Durchschnitt ausgeben — nicht exakt, aber für die Datenprofilierung oft ausreichend.
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 columnsFortschritt bei langen Läufen anzeigen
Die Verarbeitung einer mehrere Gigabyte großen Datei kann Minuten dauern. Fügen Sie eine Fortschrittsanzeige hinzu, damit Sie wissen, dass die Pipeline läuft, und die verbleibende Zeit abschätzen können. Zählen Sie die verarbeiteten Bytes oder Zeilen und vergleichen Sie sie mit der Dateigröße. Die Bibliothek tqdm macht dies mit ihrem Wrapper tqdm(reader) besonders einfach. Auch ohne tqdm liefert die Ausgabe einer Statuszeile alle 10 Blöcke während langer Batch-Jobs wertvolles Feedback.
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')Vor der Aggregation filtern
Wenden Sie Filter innerhalb jedes Blocks vor der Aggregation an, um das Ansammeln nicht benötigter Daten zu vermeiden. Wenn Sie beispielsweise nur Bestellungen aus dem Jahr 2024 benötigen, filtern Sie die Datumsspalte des Blocks vor der groupby-Aggregation. Dadurch verringert sich der für die Teilergebnisse benötigte Speicher, und der abschließende Verkettungsschritt wird schneller. Verschieben Sie Filter stets so weit wie möglich an den Anfang der Pipeline — dies ist ein grundlegendes Prinzip effizienter Datenverarbeitung.
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)Zwischenergebnisse speichern
Bei sehr lange laufenden Aufgaben sollten Sie Zwischenergebnisse regelmäßig speichern, damit Sie den Vorgang bei einer Unterbrechung ab einem Checkpoint fortsetzen können. Schreiben Sie nach jeweils N Chunks die Aggregationen pro Chunk in eine Parquet- oder CSV-Datei. Falls die Aufgabe beim Chunk 800 von 1000 fehlschlägt, können Sie die gespeicherten Aggregationen laden und dort weitermachen, wo Sie aufgehört haben, anstatt die gesamte Datei erneut zu verarbeiten. Dieses Muster für Ausfallsicherheit ist in produktiven Datenpipelines unverzichtbar.
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}')Schnelltest
Testen Sie Ihr Verständnis der Konzepte zur Datenanalyse aus dieser Lektion.
Zusammenfassung der Lektion
In dieser Lektion haben Sie Folgendes gelernt: laufende Akkumulatoren (Summe, Anzahl, Minimum/Maximum, Häufigkeits-Series) ermöglichen eine Aggregation mit konstantem Speicherbedarf über große Dateien hinweg, zweistufiges groupby (partielles groupby pro Chunk, anschließend Verketten und erneutes Gruppieren) verarbeitet gruppenübergreifende Gruppen korrekt, und frühes Filtern innerhalb jedes Chunks reduziert die Kosten des Akkumulationsschritts. Als Nächstes sehen wir uns Dask DataFrames als parallel arbeitenden, direkt einsetzbaren Ersatz für Pandas bei großen Datensätzen an.
Häufig gestellte Fragen
Ist die Lektion „Inkrementelle Aggregation über Blöcke“ kostenlos?
Ja — der vollständige Text von „Inkrementelle Aggregation über Blöcke“ ist hier im Web kostenlos zu lesen. Um sie interaktiv zu üben (integrierter Code-Editor und 24/7 KI-Tutor) und den Rest des Pandas & NumPy Academy-Kurses freizuschalten, upgrade auf CoddyKit PRO. Der Pandas & NumPy Academy-Kurs umfasst insgesamt 4 Lektionen.
Was lerne ich in „Inkrementelle Aggregation über Blöcke“?
Akkumulieren Sie laufende Anzahlen, Summen sowie Minima und Maxima über mehrere Blöcke, ohne die vollständige Datei im Speicher abzulegen. Du übst Pandas & NumPy Academy mit praktischem Code, den du direkt im Browser ausführst, und ein 24/7 KI-Tutor beantwortet deine Fragen während du die Lektion bearbeitest.
Brauche ich Erfahrung, um Pandas & NumPy Academy zu starten?
Keine Vorkenntnisse erforderlich. Pandas & NumPy Academy auf CoddyKit ist für Anfänger bis fortgeschrittene Lernende strukturiert, sodass du hier starten oder von Anfang an beginnen und in deinem eigenen Tempo voranschreiten kannst. Dies ist Lektion 2 von 4.
Wie lange dauert die Lektion „Inkrementelle Aggregation über Blöcke“?
Die meisten CoddyKit-Lektionen dauern etwa 5–10 Minuten. Jede ist kompakt und interaktiv, sodass du stetig Fortschritte machst und genau dort weitermachst, wo du aufgehört hast – im Web und in der App.
Kann ich in dieser Pandas & NumPy Academy-Lektion Code schreiben und ausführen?
Ja. Jede Pandas & NumPy Academy-Lektion enthält einen integrierten Code-Editor, sodass du echten Code direkt in deinem Browser schreibst und ausführst und sofort KI-Feedback erhältst — ohne lokale Einrichtung erforderlich.
Alle Lektionen in diesem Kurs
- CSV-Streaming mit chunksize
- Inkrementelle Aggregation über Blöcke
- Einführung in Dask DataFrames
- Parquet: Schneller spaltenbasierter Speicher