Cześć! Uczymy ludzi pracy z dużymi danymi. Nie da się wyobrazić programu edukacyjnego dotyczącego dużych danych bez własnego klastra, na którym wszyscy uczestnicy wspólnie pracują. Z tego powodu w naszym programie zawsze taki klaster istnieje 🙂 Zajmujemy się jego konfiguracją, tuningiem i administracją, a uczestnicy uruchamiają tam zadania MapReduce i korzystają z Sparka.
W tym poście opowiemy, jak rozwiązaliśmy problem nierównomiernego obciążenia klastra, pisząc własny autoskalator, używając chmury. .
Problem
Nasz klaster nie jest używany w typowy sposób. Wykorzystanie jest bardzo nierównomierne. Na przykład, są praktyczne zajęcia, kiedy wszyscy 30 uczestników i wykładowca wchodzą na klaster i zaczynają z niego korzystać. Są też dni przed ostatecznymi terminami, kiedy obciążenie znacznie rośnie. W pozostałym czasie klaster działa w trybie niedociążenia.
Rozwiązanie nr 1 – to utrzymywanie klastra, który wytrzyma szczytowe obciążenia, ale będzie nieczynny w pozostałym czasie.
Rozwiązanie nr 2 – to utrzymywanie małego klastra, do którego ręcznie dodajemy węzły przed zajęciami i podczas szczytowych obciążeń.
Rozwiązanie nr 3 – to utrzymywanie małego klastra i napisanie autoskalatora, który będzie monitorował bieżące obciążenie klastra i sam, używając różnych API, dodawał i usuwał węzły z klastra.
W tym poście będziemy mówić o rozwiązaniu nr 3. Taki autoskalator bardzo zależy od czynników zewnętrznych, a nie wewnętrznych, i dostawcy często go nie udostępniają. Korzystamy z infrastruktury chmurowej Mail.ru Cloud Solutions i napisaliśmy autoskalator, używając API MCS. A ponieważ uczymy pracy z danymi, postanowiliśmy pokazać, jak możesz napisać podobny autoskalator dla swoich potrzeb i używać go z własną chmurą.
Wymagania wstępne
Po pierwsze, musisz mieć klaster Hadoop. My na przykład korzystamy z dystrybucji HDP.
Aby węzły mogły szybko się dodawać i usuwać, musisz mieć określony podział ról w węzłach.
- Węzeł master. Tu nie ma nic do wyjaśnienia: główny węzeł klastra, na którym uruchamiany jest na przykład sterownik Sparka, jeśli korzystasz z trybu interaktywnego.
- Węzeł danych. To węzeł, na którym przechowywane są dane na HDFS i na którym odbywają się obliczenia.
- Węzeł obliczeniowy. To węzeł, na którym nic nie jest przechowywane w HDFS, ale odbywają się na nim obliczenia.
Ważna kwestia. Autoskalowanie będzie odbywać się za pomocą węzłów trzeciego typu. Jeśli zaczniesz odejmować i dodawać węzły drugiego typu, to czas reakcji znacznie się obniży - dekomisja i rekonfiguracja zajmą godziny w Twoim klastrze. Oczywiście, to nie jest to, czego oczekujesz od autoskalowania. To znaczy, że węzły pierwszego i drugiego typu pozostawiamy bez zmian. Będą one stanowić minimalnie działający klaster, który będzie istniał przez cały czas działania programu.
Zatem nasz autoskalator został napisany w Pythonie 3, korzysta z API Ambari do zarządzania usługami klastra, używa (MCS) do uruchamiania i zatrzymywania maszyn.
Architektura rozwiązania
- Moduł
autoscaler.py. Zawiera trzy klasy: 1) funkcje do pracy z Ambari, 2) funkcje do pracy z MCS, 3) funkcje związane bezpośrednio z logiką działania autoskalatora. - Skrypt
observer.py. W zasadzie składa się z różnych reguł: kiedy i w jakich momentach wywoływać funkcje autoskalatora. - Plik z parametrami konfiguracyjnymi
config.py. Zawiera na przykład listę węzłów dozwolonych do autoskalowania oraz inne parametry wpływające na to, jak długo czekać od momentu dodania nowego węzła. Znajdują się tam również znaczniki czasowe rozpoczęcia zajęć, aby maksymalna dozwolona konfiguracja klastra została uruchomiona przed zajęciami.
Przyjrzyjmy się teraz fragmentom kodu zawartym w dwóch pierwszych plikach.
1. Moduł autoscaler.py
Klasa Ambari
Oto fragment kodu zawierający klasę Ambari:
klasa 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":"Zatrzymaj wszystkie komponenty hosta",
"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 = 'Zgłoszenie przyjęte'
else:
message = req.status_code
return messagePowyżej znajduje się przykład realizacji funkcji stop_all_services, która zatrzymuje wszystkie usługi na odpowiednim węźle klastra.
Do klasy Ambari przekazujesz:
ambari_url, na przykład w postaci'http://localhost:8080/api/v1/clusters/',cluster_name– nazwa twojego klastra w Ambari,headers = {'X-Requested-By': 'ambari'}- a w środku
authznajduje się twój login i hasło do Ambari:auth = ('login', 'password').
Funkcja sama w sobie to nie więcej niż kilka wywołań przez REST API do Ambari. Z punktu widzenia logiki najpierw uzyskujemy listę uruchomionych usług na węźle, a następnie prosimy, aby na tym klastrze, na tym węźle, zmienić usługi z listy w stan INSTALLED. Funkcje uruchamiania wszystkich usług, zmiany stanu węzłów na Maintenance i inne wyglądają podobnie – to po prostu kilka zapytań przez API.
Klasa Mcs
Oto fragment kodu zawierający klasę Mcs:
klasa 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_codeDo klasy Mcs przekazujemy id projektu w chmurze oraz id użytkownika, a także jego hasło. W funkcji vm_turn_on Chcemy włączyć jedną z maszyn. Logika jest tu nieco bardziej skomplikowana. Na początku kodu następuje wywołanie trzech innych funkcji: 1) musimy uzyskać token, 2) musimy przekonwertować hostname na nazwę maszyny w MCS, 3) uzyskać id tej maszyny. Następnie wykonujemy po prostu zapytanie POST i uruchamiamy tę maszynę.
Oto funkcja do uzyskiwania tokena:
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.tokenKlasa Autoscaler
W tej klasie znajdują się funkcje związane z logiką działania.
Oto fragment kodu tej klasy:
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 messageNa wejściu przyjmujemy klasy Ambari i Mcs, listę węzłów, które są dozwolone do skalowania, a także parametry konfiguracji węzłów: pamięć i CPU przydzielone do węzła w YARN. Mamy także dwa wewnętrzne parametry q_ram, q_cpu, będące kolejkami. Dzięki nim przechowujemy wartości bieżącego obciążenia klastra. Jeśli widzimy, że przez ostatnie 5 minut obciążenie było stabilnie podwyższone, postanawiamy dodać +1 węzeł do klastra. To samo odnosi się do stanu niedoboru obciążenia klastra.
W powyższym kodzie przedstawiono przykład funkcji, która usuwa maszynę z klastra i zatrzymuje ją w chmurze. Najpierw następuje dekomisja YARN Nodemanager, następnie włącza się tryb Maintenance, potem zatrzymujemy wszystkie usługi na maszynie i wyłączamy wirtualną maszynę w chmurze.
2. Skrypt observer.py
Przykład kodu stamtąd:
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} został pomyślnie zwiększony".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)W nim sprawdzamy, czy spełnione są warunki do zwiększenia mocy klastra i czy są w ogóle dostępne maszyny rezerwowe, uzyskujemy nazwę jednej z nich, dodajemy ją do klastra i publikujemy o tym wiadomość w Slacku naszej drużyny. Po czym uruchamiany jest cooldown_period, kiedy nic nie dodajemy ani nie usuwamy z klastra, a jedynie monitorujemy obciążenie. Jeśli ustabilizowało się i znajduje się w korytarzu optymalnych wartości obciążenia, to po prostu kontynuujemy monitoring. Jeśli jednak jedna węzeł się nie wystarczył, to dodajemy kolejny.
Na wypadek, gdy przed nami jest zajęcie, wiemy na pewno, że jedna węzeł nie wystarczy, dlatego od razu uruchamiamy wszystkie dostępne węzły i utrzymujemy je aktywne do końca zajęcia. Dzieje się to za pomocą listy znaczników czasowych zajęć.
Podsumowanie
Autoskaler to dobre i wygodne rozwiązanie w przypadku, gdy obserwujesz nierównomierne obciążenie klastra. Jednocześnie osiągasz pożądaną konfigurację klastra na wysokie obciążenia i nie utrzymujesz tego klastra w czasie niedoboru obciążenia, oszczędzając środki. A to wszystko odbywa się automatycznie bez twojego udziału. Sam autoskaler to nic innego jak zestaw zapytań do API menedżera klastra i API dostawcy chmurowego, zapisanych według określonej logiki. O czym należy pamiętać – to o podziale węzłów na 3 typy, jak już wcześniej pisaliśmy. I będziecie szczęśliwi.
Źródło: habr.com
