CBDM 6 MapReduce-2 and Hadoop
Prev: 230508-1437 CBDM 5 - MapReduce
MapReduce2
- Key is to write the map part and reduce part, the distribution happens by magic by the thing
- I generate key/value pairs in a way that keys that should be grouped together are the same or similar
- Reduce - I write as if I get
keys>[list of values] - Magic inbetween, distribution into nodes in a reasonable way etc. not by me.
Termhäufigkeit MR example
- Sl.15 - Termhäufigkeit
- Generates
(term, 1)for each time this term was found - My question: can we do
(term, 4)etc. and sum the values instead of occurrences? - A: it’s a trade-off: more pairs = more networktransfer, but code is easier and prolly faster:
- Generates
- Sl. 21 - Vektormodell doc similarity through MapReduce
- Trick: one can put information inside keys if one can, such as:
(d1,d3,6*6): we put the product of vector lengths of d1/d3 inside the key (we can take the root later)
- Trick: one can put information inside keys if one can, such as:
PageRank + MapReduce example
PR
-
PR of a page is based on NUMBER and PR of websites linking to it
- Sl.23
-
Google Ranking ~= $relevanceToSearchTerms*PageRankOfPage$
-
Modelled as a person clicking links and the davon ausgehende probabilities
- A page linking t omany makes its links be ‘worth’ less
-
Formula:
-
- for each incoming link, sum (its PR / by the number of ausgehende links of that page)
- all this with a Dampfungfaktor $d$
-
-
Two ways to calculate:
- Linear Gleichung with many unknowns
- Recursive, therefore we start with some values like 1 and then iterate, each iteration changing teh PR of all pages:
- (Columns are the iteration steps)
PR + MR (iterative option)
- Paralellisierbar, weil
- PR of pages independent at each step
- Next step depends only on prev step
- Process: Sl.27
- Trick: same input shape as output, to iterate
Clusteranalysis (K-Means) + MR
-
Cluster analysis:
- build things into clusters, so that items in clusters are as similar as possible while things in different clusters should be as little similar as possible
- similarity however defined, e.g. Euclid distance
-
K-Means:
-
- In: $k$ clusters to build
- Init (once): Pick random centres for each
- Objects ‘belong’ to the clusters whose center is closest to them
- Re-calculate cluster centers
- Repeat until centers stop moving
-
-
Paralelisierbar, weil:
- objects-to-cluster Zuordnung independent from other objects
- Cluster-centers of clusters independent from other clusters
-
You can provide context - will be seen in praktikum
-
MR process:
- Map:
- given:
- existing objects (identical all iterations, given through context)
- list of all actual cluster centres
- objects-cluster mapping
- given:
- reduce:
- object with the same cluster-center
- calculation of new cluster-center for them
-
- Map:
Hadoop
Basics
-
Hadoop is an Open-source Alternative to Google’s proprietary/unavailable 2004 MR
-
Unix-ähnlich OS, Java 8/11
-
Große community etc.
-
Implements MR, no need to write OpenMPI and parallelisierungs code etc.
-
Ecosystem:
-
- “spark is better because everything in RAM, but origial MR is basis for it”
-
-
Java example of map and reduce functions Sl.36
Hadoop: Ausführungsmodell
- Tasktracker (TT) starts a set statically configured number of Map-slots and Reduce-slosts
- 10 map/reduce slots means max 10 map/reduce operations doable at the same time
- Jobtracker (JT) tracks done/failed/running tasks
- TODO: Sl. 38/39 - understand :
-
-
- Something about data on disk and RAM etc.
- Internal shuffle-sort?
- ???
-
Nel mezzo del deserto posso dire tutto quello che voglio.
comments powered by Disqus