Salut! Îi învățăm pe oameni să lucreze cu date mari. Nu putem imagina un program educațional în domeniul datelor mari fără un cluster propriu, pe care toți participanții să colaboreze. Din acest motiv, în programul nostru acesta este întotdeauna disponibil 🙂 Ne ocupăm de configurarea, optimizarea și administrarea acestuia, iar participanții rulează acolo joburi MapReduce și folosesc Spark.
În această postare, vom vorbi despre cum am rezolvat problema încărcării inegale a clusterului, scriind propriul nostru autoscalator folosind cloud-ul. .
Problema
Clusterul nostru este utilizat într-un mod care nu este tocmai tipic. Gradul de utilizare este foarte inegal. De exemplu, există sesiuni practice când toată echipa de 30 de persoane și instructorul accesează clusterul și încep să îl folosească. Sau, de asemenea, există zile înainte de termenul limită când încărcarea crește semnificativ. În restul timpului, clusterul funcționează în regim de subutilizare.
Soluția nr. 1 – este să avem un cluster care să suporte încărcările maxime, dar care să stea nefolosit în restul timpului.
Soluția nr. 2 – este să avem un cluster mic, în care să adăugăm manual noduri înainte de sesiuni și în timpul unor încărcări de vârf.
Soluția nr. 3 – este să avem un cluster mic și să scriem un autoscalator care să monitorizeze încărcarea curentă a clusterului și să adauge sau să elimine noduri din cluster, folosind diverse API-uri.
În această postare, vom discuta despre soluția nr. 3. Un astfel de autoscalator depinde foarte mult de factori externi, nu de cei interni, iar furnizorii de obicei nu îl oferă. Noi folosim infrastructura cloud Mail.ru Cloud Solutions și am scris un autoscalator folosind API-ul MCS. Și cum noi învățăm cum să lucrăm cu datele, am decis să arătăm cum puteți scrie un astfel de autoscalator pentru propriile nevoi și să-l utilizați cu cloud-ul vostru.
Cerințe preliminare
În primul rând, trebuie să aveți un cluster Hadoop. Noi, de exemplu, folosim distribuția HDP.
Pentru ca nodurile să poată fi adăugate și eliminate rapid, trebuie să aveți un anumit tip de distribuție a rolurilor pe noduri.
- Nodul master. Aici nu trebuie multe explicații: este nodul principal al clusterului, pe care rulează, de exemplu, driverul Spark, dacă folosiți modul interactiv.
- Nodul de date. Acesta este nodul pe care sunt stocate datele în HDFS și pe care se efectuează și calculele.
- Nodă de calcul. Aceasta este nodul pe care nu aveți nimic stocat pe HDFS, dar se desfășoară calcule.
Punct important. Auto-scalarea se va realiza prin noduri de tipul trei. Dacă începeți să adăugați și să eliminați noduri de tipul doi, viteza de reacție va fi semnificativ scăzută – de-comisionarea și re-comisionarea vor dura ore în clusterul dumneavoastră. Aceasta, desigur, nu este ceea ce așteptați de la auto-scalare. Așadar, nodurile de tipul unu și doi nu sunt atinse. Ele vor reprezenta un cluster minim viabil, care va exista pe durata programului.
Așadar, auto-scalerul nostru este scris în Python 3, folosește API-ul Ambari pentru gestionarea serviciilor din cluster, utilizează (MCS) pentru a porni și opri mașinile.
Arhitectura soluției
- Modulul
autoscaler.py. Acesta conține trei clase: 1) funcții pentru lucrul cu Ambari, 2) funcții pentru lucrul cu MCS, 3) funcții direct legate de logica funcționării auto-scalerului. - Script
observer.py. Practic, constă în diferite reguli: când și în ce momente să apelăm funcțiile auto-scalerului. - Fișierul cu parametrii de configurare
config.py. Acolo se află, de exemplu, lista nodurilor permise pentru auto-scalare și alți parametri care influențează, de exemplu, cât timp să așteptați de la momentul adăugării unui nod nou. Acolo se află, de asemenea, timpii de început pentru sesiuni, astfel încât, înainte de sesiune, să fie activată configurația maximă permisă a cluster-ului.
Să ne uităm acum la bucățile de cod care se află în primele două fișiere.
1. Modulul autoscaler.py
Clasa Ambari
Așa arată o bucată de cod care conține clasa Ambari:
class Ambari:
def __init__(self, ambari_url, cluster_name, headers, auth):
self.ambari_url = ambari_url
self.cluster_name = cluster_name
self.headers = headers
self.auth = auth
def stop_all_services(self, hostname):
url = self.ambari_url + self.cluster_name + '/hosts/' + hostname + '/host_components/'
url2 = self.ambari_url + self.cluster_name + '/hosts/' + hostname
req0 = requests.get(url2, headers=self.headers, auth=self.auth)
services = req0.json()['host_components']
services_list = list(map(lambda x: x['HostRoles']['component_name'], services))
data = {
"RequestInfo": {
"context":"Opriți toate componentele gazdelor",
"operation_level": {
"level":"HOST",
"cluster_name": self.cluster_name,
"host_names": hostname
},
"query":"HostRoles/component_name.in({0})".format(",".join(services_list))
},
"Body": {
"HostRoles": {
"state":"INSTALLED"
}
}
}
req = requests.put(url, data=json.dumps(data), headers=self.headers, auth=self.auth)
if req.status_code in [200, 201, 202]:
message = 'Cererea a fost acceptată'
else:
message = req.status_code
return messagePentru exemplificare, se poate consulta implementarea funcției stop_all_services, care oprește toate serviciile pe nodul dorit al clusterei.
În această clasă Ambari transmiteți:
ambari_url, de exemplu, de forma'http://localhost:8080/api/v1/clustere/',cluster_name– numele cluster-ului dumneavoastră în Ambari,headers = {'X-Requested-By': 'ambari'}- iar în interior
authse află numele de utilizator și parola dumneavoastră de la Ambari:auth = ('login', 'password').
Funcția în sine reprezintă nu mai mult decât câteva apeluri prin REST API la Ambari. Din punct de vedere logic, mai întâi obținem lista serviciilor active pe nod, apoi cerem la acest cluster, pe acest nod, să mutăm serviciile din listă în starea INSTALLED. Funcțiile pentru a porni toate serviciile, pentru a muta nodurile în starea Maintenance și altele, arată asemănător – sunt pur și simplu câteva cereri prin API.
Clasa Mcs
Așa arată o bucată de cod care conține clasa Mcs:
class Mcs:
def __init__(self, id1, id2, password):
self.id1 = id1
self.id2 = id2
self.password = password
self.mcs_host = 'https://infra.mail.ru:8774/v2.1'
def vm_turn_on(self, hostname):
self.token = self.get_mcs_token()
host = self.hostname_to_vmname(hostname)
vm_id = self.get_vm_id(host)
mcs_url1 = self.mcs_host + '/servers/' + self.vm_id + '/action'
headers = {
'X-Auth-Token': '{0}'.format(self.token),
'Content-Type': 'application/json'
}
data = {'os-start' : 'null'}
mcs = requests.post(mcs_url1, data=json.dumps(data), headers=headers)
return mcs.status_codeÎn această clasă Mcs transmitem id-ul proiectului în cloud și id-ul utilizatorului, precum și parola acestuia. În funcția vm_turn_on Vrem să activăm una dintre mașini. Logica aici este puțin mai complexă. La începutul codului, există apeluri către trei alte funcții: 1) trebuie să obținem un token, 2) trebuie să convertim hostname-ul în numele mașinii din MCS, 3) să obținem id-ul acestei mașini. Apoi facem pur și simplu o cerere POST și lansăm această mașină.
Iată cum arată funcția pentru obținerea tokenului:
def get_mcs_token(self):
url = 'https://infra.mail.ru:35357/v3/auth/tokens?nocatalog'
headers = {'Content-Type': 'application/json'}
data = {
'auth': {
'identity': {
'methods': ['password'],
'password': {
'user': {
'id': self.id1,
'password': self.password
}
}
},
'scope': {
'project': {
'id': self.id2
}
}
}
}
params = (('nocatalog', ''),)
req = requests.post(url, data=json.dumps(data), headers=headers, params=params)
self.token = req.headers['X-Subject-Token']
return self.tokenClasa Autoscaler
În această clasă se află funcțiile care se referă la logica de funcționare propriu-zisă.
Iată un fragment de cod al acestei clase:
class Autoscaler:
def __init__(self, ambari, mcs, scaling_hosts, yarn_ram_per_node, yarn_cpu_per_node):
self.scaling_hosts = scaling_hosts
self.ambari = ambari
self.mcs = mcs
self.q_ram = deque()
self.q_cpu = deque()
self.num = 0
self.yarn_ram_per_node = yarn_ram_per_node
self.yarn_cpu_per_node = yarn_cpu_per_node
def scale_down(self, hostname):
flag1 = flag2 = flag3 = flag4 = flag5 = False
if hostname in self.scaling_hosts:
while True:
time.sleep(5)
status1 = self.ambari.decommission_nodemanager(hostname)
if status1 == 'Request accepted' or status1 == 500:
flag1 = True
logging.info('Decomission request accepted: {0}'.format(flag1))
break
while True:
time.sleep(5)
status3 = self.ambari.check_service(hostname, 'NODEMANAGER')
if status3 == 'INSTALLED':
flag3 = True
logging.info('Nodemaneger decommissioned: {0}'.format(flag3))
break
while True:
time.sleep(5)
status2 = self.ambari.maintenance_on(hostname)
if status2 == 'Request accepted' or status2 == 500:
flag2 = True
logging.info('Maintenance request accepted: {0}'.format(flag2))
break
while True:
time.sleep(5)
status4 = self.ambari.check_maintenance(hostname, 'NODEMANAGER')
if status4 == 'ON' or status4 == 'IMPLIED_FROM_HOST':
flag4 = True
self.ambari.stop_all_services(hostname)
logging.info('Maintenance is on: {0}'.format(flag4))
logging.info('Stopping services')
break
time.sleep(90)
status5 = self.mcs.vm_turn_off(hostname)
while True:
time.sleep(5)
status5 = self.mcs.get_vm_info(hostname)['server']['status']
if status5 == 'SHUTOFF':
flag5 = True
logging.info('VM is turned off: {0}'.format(flag5))
break
if flag1 and flag2 and flag3 and flag4 and flag5:
message = 'Success'
logging.info('Scale-down finished')
logging.info('Cooldown period has started. Wait for several minutes')
return messageLa intrare luăm clasele Ambari și Mcs, lista nodurilor care sunt permise pentru scalare, precum și parametrii de configurare a nodurilor: memorie și CPU alocate pentru nod în YARN. Există, de asemenea, 2 parametri interni q_ram, q_cpu, care sunt cozi. Cu ajutorul lor stocăm valorile actuale ale încărcării clusterei. Dacă observăm că în ultimele 5 minute a existat o încărcare constantă crescută, luăm decizia de a adăuga +1 nod în cluster. Același lucru este valabil și pentru starea de subîncărcare a clusterei.
În codul de mai sus este un exemplu de funcție care elimină o mașină din cluster și o oprește în cloud. La început se face decommission YARN Nodemanager, apoi se activează modul Maintenance, apoi oprim toate serviciile de pe mașină și o închidem în cloud.
2. Scriptul observer.py
Exemplu de cod de acolo:
if scaler.assert_up(config.scale_up_thresholds) == True:
hostname = cloud.get_vm_to_up(config.scaling_hosts)
if hostname != None:
status1 = scaler.scale_up(hostname)
if status1 == 'Success':
text = {"text": "{0} a fost scalat cu succes".format(hostname)}
post = {"text": "{0}".format(text)}
json_data = json.dumps(post)
req = requests.post(webhook, data=json_data.encode('ascii'), headers={'Content-Type': 'application/json'})
time.sleep(config.cooldown_period*60)În acesta, verificăm dacă condițiile pentru creșterea capacităților clusterei sunt întrunite și dacă există mașini disponibile, obținem numele gazdei uneia dintre ele, o adăugăm în cluster și publicăm un mesaj în Slack echipei noastre. Apoi se pornește cooldown_period, când nu adăugăm și nu eliminăm nimic din cluster, ci doar monitorizăm încărcarea. Dacă aceasta s-a stabilizat și se află în intervalul optim al valorilor de încărcare, continuăm monitorizarea. Dacă însă o nodă nu este suficientă, adăugăm încă una.
În cazul în care avem în față o sesiune, știm deja că o nodă nu va fi suficientă, așa că pornim imediat toate nodurile disponibile și le menținem active până la finalizarea sesiunii. Acest lucru se face folosind o listă de timpi de sesiuni.
Concluzie
Auto-scalerul este o soluție bună și convenabilă pentru cazurile când aveți o încărcare inegală a clusterei. Obțineți simultan configurația necesară a clusterei pentru sarcini de vârf și nu mențineți această clusteră în timpul sub-încărcării, economisind astfel bani. Și, de asemenea, totul se desfășoară automatizat fără participarea dumneavoastră. Auto-scalerul în sine nu este altceva decât un set de cereri către API-ul managerului de cluster și API-ul furnizorului de cloud, scrise conform unei logici specifice. Ce trebuie să rețineți cu siguranță este împărțirea nodurilor în 3 tipuri, așa cum am menționat anterior. Și fericirea va veni.
Sursa: habr.com
