Kertymäaggregointi paloissa
Kerryttäkää laskureita, summia ja minimi- sekä maksimiarvoja paloittain tallentamatta koko tiedostoa muistiin.
Kertymäaggregointi paloissa on ilmainen Pandas & NumPy Academy-oppitunti CoddyKitissä. Tämä on oppitunti 2/4. Voit lukea koko oppitunnin alta ilmaiseksi ja harjoitella sen jälkeen käytännössä selaimessa sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla. Oppitunti kuuluu Pandas & NumPy Academy-oppimispolkuun, ja edistymisesi synkronoituu verkon ja CoddyKit-sovelluksen välillä. Pandas & NumPy Academy-kurssilla on yhteensä 4 oppituntia.
Miksi asteittainen koostaminen?
Asteittainen koostaminen on avain RAM-muistia suurempien aineistojen analysointiin ilman laskennan hajauttamista useille koneille. Sen sijaan että kaikki tiedot ladattaisiin lopullisen tilaston laskemista varten, ylläpidätte kertyviä arvoja — osittaisia summia, lukumääriä sekä minimi- ja maksimiarvoja — ja päivitätte niitä jokaisen palan käsittelyn yhteydessä. Lopullinen tulos muodostetaan näistä kevyistä kertyvistä arvoista, kun koko tiedosto on käyty läpi. Tällä mallilla voidaan käsitellä teratavujen kokoisia tiedostoja yhdellä kannettavalla tietokoneella.
Rivien laskeminen ja keskiarvon laskeminen
Keskiarvon laskeminen palojen yli edellyttää kertyvän summan ja lukumäärän seuraamista erikseen. Palakohtaisten keskiarvojen keskiarvoa ei voi laskea suoraan, koska palojen koot voivat vaihdella. Oikea kaava on total_sum / total_count. Tätä mallia voidaan soveltaa kaikkiin suureisiin, jotka voidaan hajottaa osiin: varianssilla, korrelaatiolla ja histogrammeilla on kaikilla asteittaiset laskentakaavat.
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}')Asteittainen minimin ja maksimin laskeminen
Globaalin minimi- ja maksimiarvon seuraaminen palojen yli on suoraviivaista: alustakaa arvot Pythonin float('inf')- ja float('-inf')-arvoilla ja päivittäkää ne kunkin palan minimi- ja maksimiarvoilla. Näin vältytään välitulosten tallentamiselta. Sama malli yleistyy ryhmäkohtaiseen minimiin ja maksimiin, kun käytätte ryhmätunnisteen mukaan avaintettua sanakirjaa.
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}')Asteittaiset frekvenssilaskennat
Luokittelevia sarakkeita varten ylläpitäkää kertyvää frekvenssisanakirjaa lisäämällä kunkin palan value_counts()-tulos Pandas Series -kertyvään arvoon. Koska Pandas Series -arvojen yhteenlasku kohdistaa arvot indeksitunnisteiden perusteella, myöhemmissä paloissa esiintyvät uudet luokat sisällytetään automaattisesti. Kun kaikki palat on käsitelty, lajitelkaa tulokset lukumäärän mukaan nähdäksenne koko aineiston yleisimmät luokat.
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))Asteittainen GroupBy-kooste
Kun haluatte laskea groupby-summan tai -lukumäärän palojen yli, käyttäkää groupby().agg()-operaatiota kussakin palassa ja tallentakaa tuloksena oleva Series tai DataFrame. Yhdistäkää kaikki osittaiset tulokset silmukan jälkeen ja tehkää niille toinen groupby-operaatio niiden yhdistämiseksi. Tämä kaksivaiheinen lähestymistapa käsittelee oikein ryhmät, jotka esiintyvät useissa paloissa. Näin käy usein, kun data on järjestetty päivämäärän eikä ryhmän mukaan.
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))Varianssin asteittainen laskeminen (Welfordin menetelmä)
Varianssin laskeminen palojen yli on keskiarvon laskemista haastavampaa. Naiivi kaava E[X²] - E[X]² kärsii katastrofaalisesta kumoutumisesta, kun keskiarvot ovat suuria. Welfordin online-algoritmi ylläpitää kertyvää keskiarvoa ja neliöityjen poikkeamien summaa sekä päivittää niitä jokaisen uuden arvon yhteydessä numeerisesti vakaalla tavalla. Vaikka SciPy toteuttaa tämän, toimintamallin ymmärtäminen auttaa laajentamaan sitä painotettuun varianssiin ja kovarianssiin.
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}')Asteittaisen histogrammin muodostaminen
RAM-muistia suureen tiedostoon tallennetun sarakkeen jakauman laskeminen edellyttää asteittaista histogrammia. Määrittäkää luokkavälien rajat etukäteen (pienen otoksen tai toimialatuntemuksen perusteella), käyttäkää sitten np.histogram(chunk_values, bins=edges)-kutsua kussakin palassa ja kerryttäkää lukumäärät. Lopuksi esittäkää yhdistetyt lukumäärät pylväskaaviona. Näin suoratoistojärjestelmät, kuten Kafka Streams ja Flink, laskevat likimääräisiä histogrammeja.
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], '...')Yksilöllisten arvojen likimääräinen seuranta
Tarkkojen erillisten arvojen laskeminen palojen yli edellyttää kaikkien yksilöllisten arvojen tallentamista — niitä voi olla miljoonia. Käyttäkää suuren mittakaavan likimääräisiin laskentoihin HyperLogLog-tietorakennetta, joka on saatavilla Pythonissa hyperloglog-kirjaston kautta. Vaihtoehtoisesti voitte seurata kunkin palan yksilöllisiä arvoja joukon avulla ja muodostaa niiden unionin, mutta joukko kasvaa rajatta. Edullista likiarvoa varten käyttäkää pd.Series.nunique()-funktiota kussakin palassa ja ilmoittakaa keskiarvo — tulos ei ole tarkka, mutta riittää usein datan profilointiin.
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 columnsEdistymisen ilmoittaminen pitkien ajojen aikana
Usean gigatavun tiedoston käsittely voi kestää minuutteja. Lisätkää edistymisen ilmoittaminen, jotta tiedätte putkilinjan olevan käynnissä ja voitte arvioida jäljellä olevan ajan. Laskekaa käsiteltyjen tavujen tai rivien määrä ja verratkaa sitä tiedoston kokoon. tqdm-kirjasto tekee tästä helppoa käyttämällä tqdm(reader)-käärettä. Ilman tqdm-kirjastoakin tilarivin tulostaminen 10 palan välein antaa arvokasta palautetta pitkien eräajojen aikana.
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')Suodattaminen ennen koostamista
Käyttäkää suodattimia kunkin palan sisällä ennen koostamista, jotta ei-toivottua dataa ei kerrytetä. Jos olette esimerkiksi kiinnostuneita vain vuoden 2024 tilauksista, suodattakaa palan päivämääräsarake ennen groupby-operaatiota. Tämä vähentää osittaisten tulosten tarvitsemaa muistia ja nopeuttaa lopullista yhdistämisvaihetta. Siirtäkää suodattimet putkilinjassa aina mahdollisimman aikaisin — tämä on tehokkaan datankäsittelyn perusperiaate.
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)Välitulosten tallentaminen
Hyvin pitkään kestävissä tehtävissä välitulokset kannattaa tallentaa säännöllisesti, jotta voitte jatkaa tarkistuspisteestä, jos prosessi keskeytyy. Kirjoittakaa kunkin lohkon aggregaatit Parquet- tai CSV-tiedostoon aina N lohkon käsittelyn jälkeen. Jos tehtävä epäonnistuu lohkossa 800, kun lohkoja on yhteensä 1000, voitte ladata tallennetut aggregaatit ja jatkaa siitä, mihin jäitte, sen sijaan että käsittelisitte koko tiedoston uudelleen. Tämä vikasietoisuusmalli on olennainen tuotantokäytössä olevissa datan käsittelyputkissa.
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}')Pikatarkistus
Testatkaa, kuinka hyvin ymmärrätte tämän oppitunnin data-analyysin käsitteet.
Oppitunnin yhteenveto
Tällä oppitunnilla opitte, että kertyvät laskurit (summa, lukumäärä, minimi/maksimi ja frekvenssi-Series) mahdollistavat suurten tiedostojen vakiomuistisen aggregoinnin, kaksivaiheinen groupby (osittainen groupby kullekin lohkolle, minkä jälkeen yhdistäminen ja uusi ryhmittely) käsittelee lohkojen yli ulottuvat ryhmät oikein ja suodatus mahdollisimman aikaisin kunkin lohkon sisällä pienentää aggregointivaiheen kustannuksia. Seuraavaksi tutustumme Dask DataFrameihin, jotka korvaavat Pandasin rinnakkaisesti suurten aineistojen käsittelyssä.
Opi Python tekoälytuutorin avulla — ilmaiseksi
Kirjoita ja suorita oikeaa koodia selaimessa, saa välitöntä apua tekoälytuutorilta ympäri vuorokauden ja jatka siitä, mihin jäit, verkossa tai sovelluksessa.
- Kurssit
- 30
- Oppitunnit
- 120
Usein kysytyt kysymykset
Onko oppitunti ”Kertymäaggregointi paloissa” ilmainen?
Kyllä – oppitunnin ”Kertymäaggregointi paloissa” koko tekstin voi lukea täällä verkossa ilmaiseksi. Jos haluat harjoitella interaktiivisesti sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla sekä avata koko Pandas & NumPy Academy-kurssin, päivitä CoddyKit PROhon. Pandas & NumPy Academy-kurssilla on yhteensä 4 oppituntia.
Mitä opin oppitunnilla ”Kertymäaggregointi paloissa”?
Kerryttäkää laskureita, summia ja minimi- sekä maksimiarvoja paloittain tallentamatta koko tiedostoa muistiin. Harjoittelet Pandas & NumPy Academy-aihetta koodilla, jonka suoritat suoraan selaimessa. Ympäri vuorokauden käytettävissä oleva tekoälytuutori vastaa kysymyksiisi oppitunnin aikana.
Tarvitsenko kokemusta aloittaakseni Pandas & NumPy Academy-opiskelun?
Aiempi kokemus ei ole tarpeen. CoddyKitin Pandas & NumPy Academy-oppimispolku sopii vasta-alkajista edistyneisiin, joten voit aloittaa tästä tai alusta ja edetä omaan tahtiisi. Tämä on oppitunti 2/4.
Kuinka kauan ”Kertymäaggregointi paloissa”-oppitunnin suorittaminen kestää?
Useimmat CoddyKitin oppitunnit kestävät noin 5–10 minuuttia. Jokainen oppitunti on lyhyt ja interaktiivinen, joten edistyt tasaisesti ja voit jatkaa siitä, mihin jäit – sekä verkossa että sovelluksessa.
Voinko kirjoittaa ja suorittaa koodia tällä Pandas & NumPy Academy-oppitunnilla?
Kyllä. Jokainen Pandas & NumPy Academy-oppitunti sisältää sisäänrakennetun koodieditorin, joten voit kirjoittaa ja suorittaa oikeaa koodia suoraan selaimessa ja saada välitöntä palautetta tekoälyltä – paikallista asennusta ei tarvita.
Kaikki tämän kurssin oppitunnit
- CSV:n suoratoisto chunksize-parametrilla
- Kertymäaggregointi paloissa
- Johdatus Dask DataFrame -rakenteisiin
- Parquet: nopea sarakepohjainen tallennus