CPD1
CPD1
Distribuit
Introducere
Cosmina Ivan
[Link]@[Link]
[Link]@[Link]
Departamentul de Calculatoare
UNIVERSITATEA TEHNICA CLUJ NAPOCA
Conţinut (1)
Concepte introductive
Auxiliar
– Consistenţă şi replicare: modele si protocoale .
– Toleranta la erori in sisteme distribuite.
3
Organizatoric
Metode de evaluare
4
5 generaţii de arhitecturi
• prima generaţie (1956) - CPU, VF, PC, CPU implicată în toate operaţiile de acces la
memorie şi de I/O, masini cu acumulator, limbaj de asamblare (Eniac, IBM701)
• a doua generaţie (1967) - registre index, aritmetica virgula flotantă,multiplexare
memorie,procesoare de I/O, limbaje de nivel inalt, procesare batch, compilatoare,
biblioteci (IBM 7030,CDC1604)
• a treia generaţie (1978)- control microprogamat, pipelining şi memorii cache , SO
time- sharing ce folosesc memoria virtuală (IBM 360, CDC 6600 PDP8).
• a patra generaţie (1989) - calculatoare paralele si sisteme distribuite, în diverse
arhitecturi cu memorie partajată sau distribuită, VLSI, SO pentru multiprocesare,
limbaje, biblioteci şi compilatoare speciale pentru paralelism, tools-uri software
(VAX9000, CrayX-MP, IBM 4090)
• a cincea generaţie (Tflops - Pflops realizează procesare masiv-paralelă,arhitecturi
scalabile şi heterogene de calculatoare cu memorie partajată sau sisteme cu transfer
de mesaje , Java, microkernel/ multithreading, S paralele si distribuite (Fujitsu,
CrayMachine, Intel-Paragon, Dash, KSR, SGI Origin)
(K Hwang, Scalable parallel computing, 1998)
5
Paralel vs distribuit
Sisteme de calcul cu mai multe elemente de procesare
Procesare concurenta
la nivel de instructiune , pipeline in procesoare
vectoriale, multiprocesare bazata pe M partajata sau
distribuita, sisteme distribuite/Internet
Aplicatii ale calculului paralel şi distribuit
Simularea informatiei
Dinamica fluidelor, Dinamica structurala – industrie automobilistica, civil
Simulari electromagnetice - radar
Planificare sisteme de control aviatice
Modelare mediu ,Simulari integrate complexe ( climat, date seismice)
Modelare biologie/sanatate, Chimie – stare solida/tranzitii
Dinamica moleculara, astrodinmica
Modelare economica/financiara
Simulari in retea, griduri, Simulari ale fluxurilor particulelor nucleare
Prelucrari grafice, procesare de imagini
Containere informationale de stocare
Analiza statistica : sanatate/asigurari
Data mining
Acces la informatie
Procesare tranzactii online/banci asigurari
Sisteme colaborative (WWW, proprietar)
Planificare piete bursiere
Integrare Informatie
Sisteme suport de decizie
Sisteme de control in timp real
E-........ Bancar, comert,.....
Educatie – e-learning , sisteme de asistare
7
Complexitate....
Alocarea resurselor
dimensionarea colectiei de elemente de procesare,
elementelor de memorie
parametrii de performanta in accesul la resurse
Performanta si scalabilitatea
translatarea cerintelor in termeni de performanta
arhitectura ofera scalabilitate?
8
Evoluţie
Architectura de sistem
arhitecturală
•transfera oferta tehnologica in performanta si capabilitati de
procesare
•rezolva diferentele generate de heterogeneitate si cerintele de
localitate a datelor, modificabile ca urmare a scalarii sistemului
sau evolutiei tehnologice
• Evolutia arhitecturilor
11
Calcul paralel
Arhitecturile paralele :
oferă o soluţie de crestere a performanţei
sistemelor de calcul (clasic, bazat doar pe evolutia
la nivelul procesoarelor)
12
Motivatie
• “Mai repede, mai mult” Problemele de rezolvat necesita
cresterea puterii de calcul
15
Caracterizarea aplicatiilor
• Compuse din taskuri.
• In aplicatii cu grad nenul de paralelism unele taskuri ruleaza
in paralel.
• Caracteristici
– granularitate
– grad de paralelism
– nivelul de paralelism
– interdependenta datelor
16
Granularitatea unui task
Granularitatea unui task
• Dependenta de:
17
Gradul de paralelism
• Gradul de paralelism
• Nivelul de paralelism:
– la nivelul procedurii
– la nivelul taskului
– la nivelul instructiunii
– la nivelul operatiei
– la nivelul microoperatiei (microcode)
18
Modele de procese
Aplicatie seriala
Modul 1
Modul 2
Modul 1
Modul 3
Single task
Multiple tasks
Modele de procese (2)
Modul 1
20
Modele de procese (3)
Modul 1
21
Taxonomii
• Obiectivele clasificarilor ( taxonomiilor) :
23
Structura multiprocesor
Procesoare P1 P2 P3 Pn
Retea
interconectare
Memorie M1 M2 Mm
24
Structura multicomputer
Retea
interconectare
P1 Pn P2
Calculatoare
M1 Mn M2
25
Structura multi-multiprocesor
Retea interconectare
26
Taxonomia Flynn
SISD – mainframe-uri, statii de lucru, PC (Single Instruction stream -Single
Data stream), masini conventionale secventiale,ce utilizeaza diverse
tehnici pentru cresterea performantelor uniprocesorului.
27
Avantajele MIMD
• Viteza mare de prelucrare, daca prelucrarea poate fi
descompusa in fire paralele, TOATE procesoarele
prelucrand simultan
28
Cerinte - MIMD (1)
1. Planificarea procesoarelor: alocarea eficienta a taskurilor la
procesoarele din sistem intr-o maniera dinamica, pe durata
evolutiei prelucrarii
30
Taxonomia Hwang
31
Modele (1)
PVP - procesoare vectoriale
Arhitectură tip “shared Memory”, module de memorie partajate ce
oferă acces rapid la date, procesoare puternice proprietar , nu integreaza
cacheuri , doar registre vectoriale si buffere de instructiuni
Reteaua de interconectare proprietar (crossbar ) (Cray C90, T90)
SMP - multiprocesoare simetrice
sistem bazat pe procesoare echivalente conectate folosind memoria
partajată şi protocoale pentru coerenta datelor
arhitectură simetrică (tip Shared-everything ) , procesoarele partajează
resursele globale disponibile, RI mag./crossbar
rulează o singura copie a SO
Limitari – reteaua de interconectare (magistrala) greu de scalat odata
proiectata(SGI, DEC Alpha, IBMR50)
32
Modele (2)
MPP – procesoare masiv paralele
sistem de procesare paralela de dimensiuni mari, procesoare clasice
şi memoria distribuită nodurilor de procesare
arhitectură shared-nothing
scalabil ( sute de noduri de procesare)
maşină asincronă MIMD, cuplat strans , folosind retele de
interconectare de înaltă performanţă (B>> , L<<)
programul conţine procese multiple, fiecare cu spaţiul privat de
adrese, interacţiune prin transfer de mesaje
procese sincronizate (operaţii de transfer de mesaje)
in anumite sisteme e posibil un singur kernel
pentru aplicaţii complexe ce posedă grad ridicat de paralelism(Intel
Paragon)
33
Modele (3)
DSM (Distributed Shared Memory)
memoriile distribuite nodurilor de procesare devin memorie globală
partajată, creand un spaţiu unic de adrese
arhitectura ce integreaza structuri d ecoerenta a datelor tip
directoare de cacheuri
necesită suport HW/SW - pentru a implementa conceptul SSI (Single
system image) (Stanford Dash - structuri de directoare,CrayT3D- suport
HW ,Treadmarks- extensii SW )
Sistem distribuit
calculatoare independente , conectate bazat pe retele convenţionale
maşinile individuale pot fi combinaţii MPP,SMP, clustere, sisteme
individuale
posedă imagini multiple de sistem (fiecare nod rulează propriul SO)
34
Modele (4)
COW –Cluster of workstation (Digital Trucluster IBM SP2, Berkeley NOW)
fiecare nod e o staţie completă (disc, mai puţin periferie), poate reprezenta
SMP sau PC
noduri conectate folosind reţea convenţională (Ethernet, FDDI, FC, switch
ATM) sau retele proprietar
interfaţa de reţea cuplată larg magistralei I-O
fiecare nod rulează propriul sistem de operare- microkernel
sistemul de operare integreaza Sw special pentru implementarea imaginii
unice de sistem (SSI)
echilibrarea încărcării, suport pentru paralelism şi comunicaţie
nodurile lucreaza colectiv asemeni unei resurse unice integrate
masina ofera disponibilitate ridicata si performanta crescuta
avantaje de cost in implementarea de masini scalabile
35
Maşini hibride tip cluster
Cluster de statii SMPs (CLUMP ) (ASCI Red (Intel), ...
36
Modele de acces la memorie
Modelul UMA(UniformMemory Access)
•memoria fizică e uniform partajată de
toate procesoarele
adresare comun
37
Modele de acces la memorie
Modelul CC-NUMA (Cache Coherent
NUMA) - toate datele partajate sunt
echidistante accesului procesoarelor
Retea de interconectare
38
Niveluri si tipuri de paralelism
39
40
41
Calcul distribuit
Def1. (Tannenbaum) - o colectie de computere independente ce apare ca
un sistem unic, coerent utilizatorilor sai – greu de realizat
42
Sistem distribuit cu elemente mobile
Internet
Mobile
phone
Printer Laptop
Camera Host site
SD acceptiune SW
aplicatii bazate pe cooperare intre procesoare/ hosturi
Caracteristici
resurse fizice si logice multiple
distributia resurselor e transparenta
independenta si cooperare intre componente
Host - componente operationale HW şi SW
Middleware – nivel intermediar , rezolva
heterogeneitatea si distribuirea entitatilor in sistem
Resurse –abstractiuni HW/SW, resurse distribuite : fizice,
date , control
44
Argumente
• Economice
– Partajarea resurselor- BD/HW costisitoare, control remote laboratoare
speciale, etc.
– Viteza /puterea de calcul oferita este posibil mai mare chiar decat a unui
mainframe.
46
Avantaje /dezavantaje
Avantaje/popularitate
• Disponibilitatea sistemelor de calcul si retelelor de
comunicatie
• Partajarea resurselor
• Scalabilitate
• Toleranta la erori
Dezavantaje:
• Puncte de cadere multiple: caderea nodurilor/ sau legaturilor
de comunicatie.
• Aspecte de securitate : mai multe oportunitati pentru atacuri
neautorizate.
47
Proprietati
Descentralizare /distributie
48
Sistem deschis
Sistem deschis - poate interactiona cu servicii ( interfete) oferite de alte sisteme
deschise, indiferent de infrastructura, independent de heterogeneitatea
HW/SW sau de limbaj.
Caracteristici:
Sistemul este conform unor interfete predefinite, publice
Este asigurata conformitatea componentelor sale cu standarde publice
(implementat folosind componente testate pentru conformanta)
suporta portabilitatea aplicatiei, ofera interoperabilitate
poate fi usor extins/ modificat/reimplementat
Concept corelat cu cele de interoperabilitate/portabilitate
49
Heterogeneitate/flexibilitate
Heterogeneitate
generata de varietatea tehnologiilor utilizate pentru implementarea
platformelor, managementului datelor si aplicatiilor
heterogeneitatea platformelor (retelelor, SO, hw) , a limbajelor a
componentelor implementate de catre dezvoltatori diferiti
rezolvata de un nivel intermediar – middleware, masini virtuale
Flexibilitate
solutie clasica - kernel monolitic - nonflexibila , sistem de fisiere/ sistem
de directoare, management complex al proceselor,
noi solutii - microkernel - servere de nivel utilizator pentru servicii sistem,
(mecanism IPC, MM, management/planificare de procese redusa)
50
Toleranta la erori/ partajare
resurse
Toleranta la erori- sistemul va opera chiar in prezenta erorilor =>determina
modalitatea de proiectare a sistemului distribuit
solutii de management – identifica si inlocuiesc componentele ce
genereaza erori. Sisteme HW – redundanta fizica, Sisteme SW –
replicari de servere, de date, tranzactii specializate
Caracteristici
independenta de L, SO, retea
utilizeaza functionalitatea de baza oferita de
sistemele de operare existente ( RPC, NFS ,
DCOM)
protocoalele / interfetele utilizate la fiecare
nivel middleware sunt identice
Middleware -caracteristici
Servicii middleware
Tipuri de middleware
orientate tranzactii(TP)- procesare tranzactii distribuite in BD multiple
orientate mesaj (MOM)- trafic sigur de mesaje intre resursele SD
orientate apel remote de proceduri (RMI)
orientate obiectual – invocari remote de obiecte(DOM)
54
Caracteristici/ cerinte pentru sisteme distribuite
heterogeneitate ( HW, SW)/ numar de procesoare ( noduri in sistem)
concurenta accesului la diverse resurse (partajarea resurselor, retele de
interconectare stari partajate )
cooperare/comunicare, gradul de sincronizare
inexistenta unui mecanism de coordonare globala
Implementarea de mecanisme de tolerare a diverselor tipuri de caderi
in sistem
55
Scalabilitatea
Scalabilitatea – defineste modul in care sistemul suporta cresterea numarului de utilizatori,
hosturi, domenii administrative, volumului datelor, tranzactiilor efectuate, distantei intre
noduri etc. (multe sisteme suporta doar scalabilitatea dimensiunii)
Tehnici de scalare
• Distributie - partitionare date/procesare pe mai multe sisteme ( ex. Apeti Java,
DNS, WWW)
• Replicare – copii ale datelor pe diverse masini (servere de fisiere, BD replicate,
site-uri replicate)
• Cachare - procese client pot accesa copii locale ale datelor: cache web, (browser,
proxy) , cache fisiere(la client/server)
• Utilizarea de comunicatii cu predilectie asincrone
Existenta unor servicii centralizate – server unic pentru utilizatori, date centralizate,
algoritmi centralizati - Solutie – algoritmi decentralizati – nici un host nu detine starea
sistemului, hosturile iau decizii bazat pe info locale, caderea unei masini nu “ distruge”
algoritmul.
56
Modele in sisteme distribuite
• Modelul de interactiune procese : modelul client-server si variatii (servere
multiple cooperante ce ofera un serviciu, cod mobil, agenti mobili), respectiv
modelul peer-to-peer (procesele pereche coopereaza si implementeaza
mecanisme de notificare evenimente sau comunicatii de grup).
– Factori ce afecteaza interactiunea : performanta canalelor,solutiile de
comunicare alese: tranzient/persistent , sincron/ asincron
60
Modelul spaţiului de memorie
•
partajat
set de adrese de memorie partajate de procesoare
• comunicatie Implicita bazata pe Memorie
• orice procesor poate referi direct orice locatie de M
• scrierile in spatiul de adrese partajat sunt vizibile
celorlalte procese/threaduri
• operatii speciale , atomice pentru sincronizare si
comunicare
• arhitecturi cuplate strans
Limitari :
cresterea numarului de procesoare conduce la
necesitatea implementarii unori mecanisme
complexe care sa rezolve accesele concurente la
acceasi locatie de memorie
61
Modelul transferului de
mesaje…
– comunicatie bazata pe operatii I/O explicite
– usor de proiectat , mult mai scalabile
• Modelul de programare
– acces direct doar la spatiul de adrese private (M locala)
– comunicatia bazata pe mesaje explicite (send/receive)
62
Abstracţiunea Message-Passing
Match t Receive
, Y, P,t
Adresa Y
Send X, Q, t
Spatiu de
Adresa x Spatiu de adrese local
adrese local
Pr oces P Pr oces Q
63
Modele de interactiune
Modelul obiectual
Obiect = resursa partajata , poseda identitate unica, se afla in server si
ofera facilitati de acces la distanta
obiectul client poseda un proxy – implementarea interfetei obiectului
pot fi colocate in alte hosturi mentinand identitatea ( migrare obiecte)
Observatii
distinctie minimalista intre servere - clienti , reflecta modelul de
interactiune “ peer to peer” ( un obiect poate fi client/server,
solicita servicii altor obiecte)
nivelul middleware implementeaza mecanisme de localizare a
obiectelor si de comunicare/interactiune
65
Modelul minicomputer
Mini-
computer
ARPA
Mini- net Mini-
computer computer
66
Modelul Workstation
Workstation
Workstation Workstation
67
Modelul Workstation-Server
• presupune existenţa unor staţii simple
pentru procesare locală a aplicaţiilor
interactive , iar operaţii de acces la
Workstation fişiere, tipărire, http sunt transferate
unor servere specializate.
Workstation Workstation
• Serverele sunt sisteme puternice fiecare
dedicat unui anumit tip de serviciu.
100Gbps
LAN
• Modelul de comunicaţie client - server -
tip apel de procedură la distanţă sau
invocare remote de obiecte .
Mini- Mini- Mini-
Computer Computer Computer
file server http server ... server • Modalitatea de interacţiune este
realizată prin apelul procesului server de
către procesul client, fără a presupune
migrare.
68
Modelul Ferma de procesoare
terminale diskless
100Gbps
LAN
pentru fiecare utilizator va putea fi
alocat numărul necesar de procesoare.
69
Modelul Cluster
serverul conţine staţii sau sisteme
Workstation
personale conectate printr-o reţea de mare
viteză, având drept scop performanţă
Workstation Workstation oferită aplicaţiilor - pot procesa cereri în
paralel.
100Gbps
LAN
http server2
http server1 http server N
1Gbps SAN
70
Modelul grid
a apărut în intenţia de a colecta puterea
de calcul a sistemelor tip supercomputer,
Workstation
şi a clusterelor distribuite cu scopul de a
le oferi ca o resursă unitară de calcul.
Workstation Workstation
71
Modele computationale - Modelul Client-Server
File server
DNS server
result result
Server
Client
Key:
Process: Computer:
Workstation
Workstation Workstation
100Gbps
LAN
72
Structuri de servere
Service
Server
Client Replication
• Availability
• Performance
Server
Client
Workstation
Server
Workstation Workstation
100Gbps
LAN
73
Servere proxy si cache
Client Web
server
Proxy
server
Client Web
Ex. Internet Service Provider server
Workstation
Workstation Workstation
100Gbps
LAN
1Gbps SAN
74
Procese peer
Application Application
Coordination Coordination
code code
Application
Coordination
Distributed whiteboard application
code
Workstation
75
Cod mobil , agenti
Client Web
Applet code server
Web
Client Applet server
Workstation
Workstation Workstation
100Gbps
LAN
76
Clienti Thin
Workstation
Workstation Workstation
100Gbps
100Gbps
LAN
LAN
1Gbps SAN
77
“The Power of the Internet “
(Source: [Link])/2012
• DOMAIN NAMES: There are 12,844,877 unique domain names (e.g. [Link]) registered
worldwide, with 428,023 new domain names registered each week. (NetNames Statistics
12/28/1999).
• BACKBONE CAPACITY: The capacity of the Internet backbone to carry information is doubling every
100 days. (U.S. Internet Council, Apr. 2004).
• DATA TRAFFIC SURPASSING VOICE: Voice traffic is growing at 10% per year or less, while data
traffic is conservatively estimated to be growing at 125% per year, meaning voice will be less than
1% of the total traffic by 2007. (Technology Futures, Inc March 2000).
• HOST COMPUTERS: In July 1999 there were 56.2 million "host" computers supporting web pages.
In July 1997 there were 19.5 million host computers, with 3.2 million hosts in July 1994, and a
mere 80,000 in July 1989. (Internet Software Consortium – Internet Domain Survey).
• TOTAL AMOUNT OF DATA: 1,570,000,000 pages, 29,400,000,000,000 bytes of text, 353,000,000
images, and 5,880,000,000,000 bytes of image data. (The Censorware Project, Jan. 26, 1999).
• EMAIL VOLUME: Average U.S. consumer will receive 1,600 commercial email messages in 2005, up
from 40 in 1999, while non-marketing and personal correspondence will more than double from
approximately 1,750 emails per year in 1999 to almost 4,000 in 2005 (Jupiter Communications,
May 2000).
• 159 million computers in the U.S., 135 million in EU, and 116 million in Asia Pacific (as of April
2000).
• WEB HITS/DAY: U.S. web pages averaged one billion hits per day (aggregate) in October 1999.
(eMarketer/Media Metrix, Nov. 1999).
78
Google
• the Google cluster has the following stats:
359 racks
31,654 machines
63,184 CPUs
63,184 Gb of RAM
sursa [Link]/2012
79
“The network really is the computer.”
30,000 servers
17
Exemplu
2
Paralel vs distribuit
83
Concluzii
Calcul paralel Calcul distribuit
... ftp- materiale ( carti, note de curs in romana, slideuri, laburi, resurse
referate/proiecte)
85