Jak stworzyć swój własny autoskalator dla klastra

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. Mail.ru Cloud Solutions.

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.

  1. 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.
  2. Węzeł danych. To węzeł, na którym przechowywane są dane na HDFS i na którym odbywają się obliczenia.
  3. 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 API Mail.ru Cloud Solutions (MCS) do uruchamiania i zatrzymywania maszyn.

Architektura rozwiązania

  1. 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.
  2. Skrypt observer.py. W zasadzie składa się z różnych reguł: kiedy i w jakich momentach wywoływać funkcje autoskalatora.
  3. 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 message

Powyż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 auth znajduje 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_code

Do 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.token

Klasa 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 message

Na 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

Kup niezawodny hosting stron z ochroną DDoS, serwery VPS VDS 🔥 Kup niezawodny hosting stron z ochroną DDoS, serwery VPS VDS - ProHoster