Kuidas luua oma auto skaler klastrile

Tere! Me koolitame inimesi suurte andmete töötlemises. Hariduse programm suurte andmete kohta on võimatu ette kujutada ilma oma klastrita, kus kõik osalejad koos töötavad. Sel põhjusel on meie programmis see alati olemas 🙂 Me tegeleme selle seadistamise, häälestamise ja haldamisega, samas kui osalejad käivitavad seal MapReduce töökohti ja kasutavad Spark'i.

Selles postituses räägime, kuidas me lahendasime klastrite ebaühtlase koormuse probleemi, kirjutades oma auto skalerit, kasutades pilve. Mail.ru Cloud Solutions.

Probleem

Meie klassrit kasutatakse mitte täiesti tavalises režiimis. Kasutuse määr on väga ebavõrdne. Näiteks on praktilised õppused, kui kõik 30 inimest ja õpetaja siseneb klastrisse ja hakkavad seda kasutama. Samuti on mõni päev enne tähtaega, kui koormus märgatavalt tõuseb. Kõik ülejäänud aega töötab klaster alakoormuse režiimis.

Lahendus nr 1 – hoida klastrit, mis talub tipukoormusi, aga mis seisab ülejäänud aja.

Lahendus #2 on hoida väikest klastrit, kuhu käsitsi lisada sõlmi enne ülesandeid ja tipptundide ajal.

Lahendus #3 on hoida väikest klastrit ja kirjutada automaatne skaaler, mis jälgib klastrite praegust koormust ning lisab ja eemaldab sõlmi klastrist erinevate API-de abil.

Selles postituses räägime lahendusest #3. Selline automaatne skaaler sõltub tugevalt välistest teguritest, mitte sisemistest, ning teenusepakkujad ei paku seda sageli. Me kasutame Mail.ru Cloud Solutions'i pilvesüsteemi ja koostasime automaatse skaaleri, kasutades MCS API-d. Kuna me õpetame andmete töötlemist, otsustasime näidata, kuidas saate kirjutada sarnase automaatse skaaleri oma eesmärkide jaoks ja kasutada seda oma pilves.

Eeltingimused

Esiteks peab teil olema Hadoopi klaster. Näiteks kasutame HDP distributsiooni.

Selleks, et sõlmed saaksid kiiresti lisanduda ja eemalduda, peab teil olema sõlmede vahel teatud rollide jaotus.

  1. Meister-sõlm. Siin pole erilisi selgitusi vaja: klastrite peamine sõlm, kus toimub näiteks Spark'i draiveri käivitamine, kui kasutate interaktiivset režiimi.
  2. Andme sõlm. See on sõlm, kus teie andmed HDFS-is on ja kus toimub ka arvutamine.
  3. Arvutussõlm. See on sõlm, kus teil ei ole HDFS-is andmeid, kuid kus toimub arvutamine.

Oluline punkt. Autoskaleerimine toimub kolmanda tüübi sõlmade arvelt. Kui hakkate teise tüübi sõlmi eemaldama ja lisama, siis reageerimiskiirus on oluliselt madalam – dekomisioon ja re-komisioon võtavad teie klastri peal tundide kaupa aega. See pole kindlasti see, mida ootaksite autoskaleerimiselt. Seega esimest ja teist tüüpi sõlmi me ei puuduta. Need moodustavad minimaalselt elujõulise klastri, mis eksisteerib kogu programmi kestuse.

Nii et meie autoskaleerija on kirjutatud Python 3-s, kasutab Ambari API-d klastri teenuste haldamiseks, kasutab Mail.ru Cloud Solutionsi (MCS) API-d masinate käivitamiseks ja peatamiseks. (MCS) для запуска и остановки машин.

Lahenduse arhitektuur

  1. Modul autoscaler.py. See sisaldab kolme klassi: 1) funktsioonid Ambari töötamiseks, 2) funktsioonid MCS-iga töötamiseks, 3) funktsioonid, mis on otseselt seotud autoskaleerija töölogikaga.
  2. Skript observer.py. Põhimõtteliselt koosneb see erinevatest reeglitest: millal ja milliseid automaatse skaleerimise funktsioone kutsuda.
  3. Konfiguratsiooniparametrite fail config.py. Seal on näiteks loetelu nodidest, mis on lubatud automaatseks skaleerimiseks, ja teised parameetrid, mis mõjutavad näiteks seda, kui kaua oodata pärast uue noodi lisamist. Seal on ka algusajamärgid, et enne õppetundi oleks käivitatud maksimaalne lubatud klastrikonfiguratsioon.

Vaatame nüüd koodi osi, mis asuvad kahes esimeses failis.

1. Moodul autoscaler.py

Klass Ambari

Nii näeb välja koodilõik, mis sisaldab klassi 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":"Stop All Host Components",
                "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 = 'Request accepted'
        else:
            message = req.status_code
        return message

Ülaltoodud näite põhjal on võimalik tutvuda funktsiooni rakendamisega stop_all_services, mis peatab kõik teenused valitud klastrinoodil.

Klassile toimub sisendi edastamine Ambari ja need on:

  • ambari_url, näiteks kujul 'http://localhost:8080/api/v1/clusters/',
  • cluster_name – teie klastrinimi Ambari's,
  • headers = {'X-Requested-By': 'ambari'}
  • ja seespool auth on teie Ambari kasutajanimi ja parool: auth = ('login', 'password').

Funktsioon ise sisaldab vaid paar REST API kutset Ambari poole. Loogiliselt vaadatuna saame alguses nimekirja node'il käivitatud teenustest ja seejärel palume selle klastris, et node'i teenused muudetaks olekusse INSTALLED. Kõik teenuste käivitamise, node'ide olekusse viimise funktsioonid Maintenance ja muud toimivad sarnasel viisil – tegemist on lihtsalt mitme API päringuga.

Klass Mcs

Nii näeb välja koodilõik, mis sisaldab klassi 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

Klassile toimub sisendi edastamine Mcs me edastame projekti id pilves ja kasutaja id ning tema parooli. Funktsioonis vm_turn_on meie soovime aktiveerida ühe masinat. Loogika siin on veidi keerulisem. Koodi alguses kutsub see välja kolm muud funktsiooni: 1) meil on vaja saada token, 2) meil on vaja konverteerida hostname masina nimeks MCS-s, 3) saada selle masina ID. Seejärel teeme lihtsalt post-päringu ja käivitame selle masina.

Nii välja näeb funktsioon tokeni saamiseks:

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

Autoscaler klass

Selles klassis on funksioonid, mis on seotud töö loogikaga.

Nii välja näeb selle klassi koodilõik:

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('Decommission 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

Me võtame sisendina klassid Ambari ja Mcs, loendi sõlmedest, mis on lubatud skaleerimiseks, ning sõlme konfiguratsiooniparameetritest: mälu ja CPU, mis on eraldatud sõlmele YARNis. Samuti on kaks sisemist parameetrit q_ram ja q_cpu, mis on järjekorrad. Me hoiame nende abil käimasoleva klastrikoormuse väärtusi. Kui näeme, et viimase viie minuti jooksul on pidevalt olnud suurenenud koormus, siis otsustame, et klastrisse on vaja lisada +1 sõlm. Sama kehtib klastrikoormuse alakoormuse puhul.

Ülaltoodud koodis on toodud näide funktsioonist, mis eemaldab masina klastrist ja peatab selle pilves. Alguses toimub dekomisjonimine YARN Nodemanager, siis lülitatakse sisse režiim Maintenance, seejärel peatame kõik teenused masinas ja lülitame välja virtuaalse masina pilves.

2. Skript observer.py

Koodinäidis sealt:

kui scaler.assert_up(config.scale_up_thresholds) == True:
        hostname = cloud.get_vm_to_up(config.scaling_hosts)
        kui hostname != None:
            status1 = scaler.scale_up(hostname)
            kui status1 == 'Success':
                text = {"text": "{0} on edukalt suurendatud".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)

Selles kontrollime, kas klastrite suurendamise tingimused on täidetud ja kas reservis on masinad, saame ühe nende hosti nime, lisame selle klastrisse ja saadame sellest teadet meie meeskonna Slackis. Seejärel käivitame cooldown_period, mil me ei lisa ega eemalda klastrist midagi, vaid jälgime koormust. Kui see stabiliseerub ja jääb optimaalse koormuse vahemikku, jätkame lihtsalt jälgimist. Kui ühe sõlme jaoks ei piisavad, siis lisame veel ühe.

Kui meil on ees tegevus ja me teame kindlalt, et ühest sõlmest ei piisa, siis käivitame kohe kõik vabad sõlmed ja hoiame need aktiivsena kuni tegevuse lõpuni. See toimub tegevuste ajatempli nimekirja abil.

Kokkuvõte

Autoskaleerija on hea ja mugav lahendus, kui teil on klastris ebaühtlane koormus. Samal ajal saavutate vajaliku klastrikonfiguratsiooni tipukoormuste jaoks ja ei pea klastrit pidevalt hoidma, kui koormus on madal, säästes seeläbi kulusid. Ja kõik see toimub automaatselt, ilma teie osaluseta. Autoskaleerija ise ei ole rohkem kui komplekt päringuid klastrihalduri ja pilveteenuse pakkuja API-le, mis on kirjutatud kindla loogika järgi. Oluline on meeles pidada, et noded tuleb jagada kolme tüüpi, nagu me varem rääkisime. Ja see toob teile õnne.

Allikas: habr.com

Osta usaldusväärne veebihosting DDoS kaitsega, VPS VDS serverid 🔥 Osta usaldusväärne veebihosting DDoS kaitsega, VPS VDS serverid | ProHoster