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) masinate kÀivitamiseks ja peatamiseks.

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