Pandas & NumPy Academy · Lekcja

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.

Lekcja 2 z 413 kroki

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 columns

Raportowanie 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.

Bezpłatny start

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

  1. Strumieniowy odczyt CSV z chunksize
  2. Przyrostowa agregacja między porcjami
  3. Wprowadzenie do obiektów Dask DataFrame
  4. Parquet: szybki magazyn kolumnowy
← Powrót do Pandas & NumPy Academy