Komplett pensumoversikt for store, distribuerte datamengder ved NTNU — med forklaringer, sentrale begreper, eksamenstips og vanlige fallgruver. Eksamensoptimalisert basert på tidligere eksamener.
Denne studieguiden dekker TDT4225 Store, distribuerte datamengder ved NTNU (7,5 stp, masternivå). Faget handler om hvordan datamengder som er for store for arbeidslageret faktisk blir lagret, indeksert, sortert og behandlet — og hvordan du regner på lagerbehov, I/O-volum og responstider i stedet for å gjette.
Slik henger innholdet sammen med vurderingen. Emnet vurderes i dag med to gruppeprosjekter (40 %) og en skriftlig skoleeksamen på 3 timer (60 %), og emnebeskrivelsen legger vekt på distribuerte systemdesign, datamodeller og spørrespråk, indeksering og lagringsmetoder, koding, replikasjon og partisjonering, transaksjoner, konsistens og konsensus, samt database-as-a-service. Eksamensarkivet vi har kalibrert oppgavene mot er fra perioden da emnet het Lagring og behandling av store datamengder (skoleeksamen 3–4 timer, hjelpemiddelkode D: ingen trykte eller håndskrevne hjelpemidler, bestemt enkel kalkulator tillatt). Kjernemetodene i de settene — hashing og blokkorganisering, flerdimensjonale indekser, ekstern sortering, I/O-dimensjonering og relasjonsalgebra med gjentatte gjennomløp — er fortsatt fagets regnehåndverk, og det er der de fleste poengene ligger på en skriftlig prøve.
Guiden er bygget rundt ni temaer. Fem av dem er eksamensforankret i arkivet 2009–2012: filstrukturer og hashing, indeksering, systemdimensjonering og I/O-volum, ekstern sortering og fletting, og relasjonsalgebra og spørreutførelse. Fire er forankret i dagens emnebeskrivelse: parallellprosessering, distribuerte systemer, NoSQL og MapReduce. Bruk de fem første til å trene regneferdighet og de fire siste til å bygge begrepsapparatet du trenger i prosjektene og i drøftingsoppgavene.
Tre ferdigheter går igjen i alt materialet:
Et gjennomgående prinsipp knytter alt sammen: flaskehalsen er transporten, ikke regningen. Nesten hver eneste vurdering i faget koker ned til å telle hvor mange ganger data må passere mellom disk og arbeidslager.
Programmeringsmodell for parallell batchbehandling av store datasett gjennom map- og reduce-faser med shuffle/sort imellom.
Når et datasett er for stort til å passe i ett arbeidslager og må prosesseres på mange maskiner samtidig, trenger vi en modell som skjuler kompleksiteten ved parallellitet, datafordeling og feiltoleranse. MapReduce (Dean & Ghemawat, Google 2004) er en slik modell: programmereren skriver bare to funksjoner, og rammeverket håndterer alt det vanskelige.
map(k1, v1) → liste(k2, v2). Kjøres uavhengig på hver inputblokk (split). Produserer mellomliggende nøkkel–verdi-par. Fordi map-oppgavene er uavhengige, kan de kjøre fullstendig parallelt.reduce(k2, liste(v2)) → liste(v3). Får alle verdier for én nøkkel samlet og aggregerer dem.Mellom map og reduce ligger det dyreste steget: shuffle. Alle mellomliggende par med samme nøkkel k2 må samles på samme reducer. Dette krever partisjonering (typisk hash(k2) mod R, der R er antall reducere), sortering på nøkkel, og overføring av data over nettet. Dette er en distribuert variant av nettopp den eksterne sorteringen og partisjoneringen vi gjør i ett enkelt system.
Inputfilen deles i splits (typisk lik blokkstørrelsen i det distribuerte filsystemet). Planleggeren prøver å kjøre en map-oppgave på den maskinen som allerede har blokken lokalt — slik unngås nettverkstrafikk. Dette er det samme prinsippet vi bruker når vi minimerer I/O-volum: flytt beregningen til dataene, ikke omvendt.
Hver oppgave er deterministisk og uten sideeffekter. Faller en maskin ut, kjøres oppgaven på nytt et annet sted fra inputblokken. Trege noder («stragglers») håndteres med spekulativ eksekvering: en kopi av oppgaven startes i parallell, og første som blir ferdig vinner.
Du skal telle forekomster av hvert produktnummer i en logg på 4 TB. Det distribuerte filsystemet bruker blokkstørrelse 128 MB.
Antall map-oppgaver: 4 TB / 128 MB = 4 194 304 MB / 128 MB = 32 768 map-oppgaver (én per split).
Map-output: hver map-funksjon sender ut par (produktnr, 1). Med en lokal combiner som forhåndsaggregerer per blokk, reduseres mellomdataene kraftig før shuffle.
Shuffle-volum: hvis det finnes 50 000 distinkte produktnumre og hver mapper i snitt ser 10 000 distinkte, blir mellomvolumet etter combiner ca. 32 768 × 10 000 par. Hver reducer (med R = 64) får da i gjennomsnitt (32 768 × 10 000)/64 par å summere. Dette illustrerer hvorfor combiner og fornuftig valg av R er avgjørende for shuffle-kostnaden.
En distribuert join i MapReduce kan gjøres enten som reduce-side join (begge tabeller emittes med koblingsnøkkel, reducer matcher) eller map-side join (den lille tabellen kringkastes til alle mappere). Dette tilsvarer valget mellom gjentatte gjennomløp og å holde den minste operanden i arbeidslager, slik vi gjør i relasjonsalgebra på én maskin.
Nøkkelformler
Vanlige feil
Eksamenstips
Laster...