Przyrostowa agregacja między porcjami
Akumuluj bieżące liczności, sumy oraz wartości min/max między porcjami bez przechowywania całego pliku w pamięci.
Przyrostowa agregacja między porcjami to bezpłatna lekcja Pandas & NumPy Academy na CoddyKit. To lekcja 2 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej Pandas & NumPy Academy, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs Pandas & NumPy Academy zawiera 4 lekcji w sumie.
Dlaczego warto stosować agregację przyrostową
Agregacja przyrostowa jest kluczem do analizowania zbiorów danych większych niż pamięć RAM bez rozdzielania obliczeń między wiele maszyn. Zamiast ładować wszystkie dane w celu obliczenia końcowej statystyki, należy utrzymywać akumulatory narastające — sumy częściowe, liczności oraz wartości minimum/maksimum — i aktualizować je przy każdym fragmencie. Po przeskanowaniu całego pliku końcowy wynik jest tworzony na podstawie tych lekkich akumulatorów. Ten wzorzec umożliwia przetwarzanie plików terabajtowych na jednym laptopie.
Zliczanie wierszy i obliczanie średniej
Obliczanie średniej dla wielu fragmentów wymaga osobnego śledzenia narastającej sumy i liczby elementów. Nie można po prostu uśrednić średnich z poszczególnych fragmentów, ponieważ fragmenty mogą mieć różne rozmiary. Prawidłowy wzór to total_sum / total_count. Wzorzec ten można zastosować do każdej wielkości, którą da się rozłożyć na składniki: wariancja, korelacja i histogramy również mają wzory przyrostowe.
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}')Przyrostowe wyznaczanie minimum i maksimum
Śledzenie globalnego minimum i maksimum w kolejnych fragmentach jest proste: należy zainicjalizować wartości za pomocą float('inf') i float('-inf'), a następnie aktualizować je na podstawie minimum/maksimum każdego fragmentu. Eliminuje to konieczność przechowywania danych pośrednich. Wzorzec można uogólnić na minimum/maksimum dla grup, utrzymując słownik, którego kluczem jest identyfikator grupy.
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}')Przyrostowe zliczanie częstości
W przypadku kolumn kategorialnych należy utrzymywać słownik częstości narastających, dodając wynik value_counts() z każdego fragmentu do akumulatora w postaci Series Pandas. Ponieważ dodawanie Series Pandas dopasowuje etykiety indeksu, nieznane kategorie pojawiające się w kolejnych fragmentach są automatycznie uwzględniane. Po przetworzeniu wszystkich fragmentów posortuj wyniki według liczby wystąpień, aby zobaczyć najczęstsze kategorie w całym zbiorze danych.
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))Przyrostowa agregacja GroupBy
Aby obliczyć sumę lub liczbę za pomocą groupby dla wielu fragmentów, zastosuj groupby().agg() w każdym fragmencie i przechowuj wynikową Series lub DataFrame. Po zakończeniu pętli połącz wszystkie wyniki częściowe i zastosuj drugie grupowanie, aby je scalić. To dwuetapowe podejście poprawnie obsługuje grupy występujące w wielu fragmentach, co jest częste, gdy dane są posortowane według daty, a nie według grupy.
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))Przyrostowe obliczanie wariancji (metoda Welforda)
Obliczanie wariancji dla wielu fragmentów jest trudniejsze niż obliczanie średniej. Naiwny wzór E[X²] - E[X]² cierpi z powodu katastrofalnego skracania dla dużych średnich. Algorytm online Welforda utrzymuje narastającą średnią oraz sumę kwadratów odchyleń i aktualizuje je przy każdej nowej wartości w sposób stabilny numerycznie. Choć SciPy implementuje ten algorytm, zrozumienie tego wzorca pozwala rozszerzyć go na wariancję ważoną i kowariancję.
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}')Tworzenie histogramu przyrostowego
Obliczenie rozkładu kolumny w pliku zbyt dużym, aby zmieścił się w pamięci RAM, wymaga użycia histogramu przyrostowego. Z góry ustal krawędzie przedziałów (na podstawie małej próbki lub wiedzy dziedzinowej), a następnie w każdym fragmencie użyj np.histogram(chunk_values, bins=edges) i gromadź liczności. Na końcu przedstaw połączone liczności na wykresie słupkowym. Tak właśnie systemy strumieniowe, takie jak Kafka Streams i Flink, obliczają przybliżone histogramy.
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], '...')Przybliżone śledzenie unikatowych wartości
Dokładne zliczanie różnych wartości w wielu fragmentach wymaga przechowywania wszystkich unikatowych wartości — potencjalnie milionów. Do przybliżonego zliczania na dużą skalę należy użyć szkicu HyperLogLog, dostępnego w Pythonie za pośrednictwem biblioteki hyperloglog. Alternatywnie można śledzić unikatowe wartości dla każdego fragmentu za pomocą zbioru i obliczyć sumę zbiorów, ale taki zbiór będzie nieograniczenie rosnąć. Aby uzyskać niedrogie przybliżenie, użyj pd.Series.nunique() dla każdego fragmentu i podaj średnią — wynik nie będzie dokładny, ale często wystarcza do profilowania danych.
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 columnsRaportowanie postępu podczas długich uruchomień
Przetwarzanie pliku o rozmiarze kilku gigabajtów może trwać kilka minut. Dodaj raportowanie postępu, aby wiedzieć, czy potok działa, i móc oszacować pozostały czas. Zliczaj przetworzone bajty lub wiersze i porównuj je z rozmiarem pliku. Biblioteka tqdm znacznie ułatwia to zadanie dzięki opakowaniu tqdm(reader). Nawet bez tqdm wyświetlanie wiersza stanu co 10 fragmentów dostarcza cennych informacji podczas długich zadań wsadowych.
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')Filtrowanie przed agregacją
Stosuj filtry w każdym fragmencie przed agregacją, aby uniknąć gromadzenia niepotrzebnych danych. Jeśli na przykład interesują Państwa tylko zamówienia z 2024 roku, należy odfiltrować odpowiednie daty w kolumnie fragmentu przed wykonaniem groupby. Ogranicza to pamięć potrzebną na wyniki częściowe i przyspiesza końcowy etap ich łączenia. Filtry należy zawsze stosować możliwie wcześnie w potoku — jest to fundamentalna zasada wydajnego przetwarzania danych.
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)Zapisywanie wyników pośrednich
W przypadku zadań wykonywanych przez bardzo długi czas należy okresowo zapisywać wyniki pośrednie, aby można było wznowić działanie od punktu kontrolnego, jeśli proces zostanie przerwany. Po każdych N fragmentach należy zapisywać agregaty dla przetworzonych fragmentów do pliku Parquet lub CSV. Jeśli zadanie zakończy się niepowodzeniem przy fragmencie 800 z 1000, można wczytać zapisane agregaty i kontynuować od przerwanego miejsca zamiast ponownie przetwarzać cały plik. Ten wzorzec odporności ma kluczowe znaczenie w produkcyjnych potokach danych.
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}')Szybkie sprawdzenie
Sprawdź swoją wiedzę na temat analizy danych z tego rozdziału.
Podsumowanie rozdziału
W tym rozdziale poznano: akumulatory działające w trakcie przetwarzania (suma, licznik, minimum/maksimum, seria częstości) umożliwiają agregowanie dużych plików przy stałym zużyciu pamięci, dwuetapowe grupowanie groupby (częściowe grupowanie dla każdego fragmentu, a następnie konkatenacja i ponowne grupowanie) prawidłowo obsługuje grupy obejmujące wiele fragmentów, a wczesne filtrowanie w obrębie każdego fragmentu zmniejsza koszt etapu agregacji. W dalszej części przyjrzymy się obiektom Dask DataFrame jako bezpośredniemu, równoległemu zamiennikowi Pandas dla dużych zbiorów danych.
Ucz się Python dzięki korepetycjom AI — za darmo
Pisz i uruchamiaj kod w przeglądarce, otrzymuj natychmiastową pomoc od korepetytora AI dostępnego 24/7 i kontynuuj naukę w sieci lub w aplikacji.
- Kursy
- 30
- Lekcje
- 120
Często zadawane pytania
Czy lekcja „Przyrostowa agregacja między porcjami” jest bezpłatna?
Tak — pełny tekst „Przyrostowa agregacja między porcjami” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu Pandas & NumPy Academy, przejdź na CoddyKit PRO. Kurs Pandas & NumPy Academy zawiera 4 lekcji w sumie.
Co nauczysz się w „Przyrostowa agregacja między porcjami”?
Akumuluj bieżące liczności, sumy oraz wartości min/max między porcjami bez przechowywania całego pliku w pamięci. Ćwiczysz Pandas & NumPy Academy z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć Pandas & NumPy Academy?
Nie wymagamy żadnego doświadczenia. Pandas & NumPy Academy w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 2 z 4.
Ile czasu zajmuje lekcja „Przyrostowa agregacja między porcjami”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji Pandas & NumPy Academy?
Tak. Każda lekcja Pandas & NumPy Academy zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Strumieniowy odczyt CSV z chunksize
- Przyrostowa agregacja między porcjami
- Wprowadzenie do obiektów Dask DataFrame
- Parquet: szybki magazyn kolumnowy