Come creare il proprio autoscalatore per un cluster

Ciao! Insegniamo alle persone a lavorare con grandi dati. Non si può immaginare un programma educativo sui big data senza il proprio cluster, dove tutti i partecipanti lavorano insieme. Per questo motivo, nel nostro programma è sempre presente 🙂 Ci occupiamo della sua configurazione, messa a punto e amministrazione, mentre i ragazzi eseguono direttamente lavori MapReduce e utilizzano Spark.

In questo post parleremo di come abbiamo risolto il problema del carico non uniforme del cluster, creando il nostro autoscalatore utilizzando il cloud Mail.ru Cloud Solutions.

Problema

Il nostro cluster non è utilizzato in modo del tutto tipico. L'utilizzo è molto irregolare. Ad esempio, ci sono sessioni pratiche in cui tutte e 30 le persone e il docente accedono al cluster e iniziano a utilizzarlo. Oppure, ci sono giorni prima della scadenza in cui il carico aumenta notevolmente. In tutto il resto del tempo, il cluster funziona in modalità sottoutilizzo.

Soluzione №1 – mantenere un cluster che possa sostenere i picchi di carico, ma che rimanga inattivo per tutto il resto del tempo.

Soluzione №2 – mantenere un piccolo cluster, aggiungendo manualmente nodi prima delle lezioni e durante i carichi di picco.

Soluzione №3 – mantenere un piccolo cluster e scrivere un autoscalatore che monitori il carico attuale del cluster e aggiunga o rimuova nodi dal cluster utilizzando varie API.

In questo post parleremo della soluzione №3. Questo autoscalatore dipende fortemente da fattori esterni piuttosto che interni, e i fornitori spesso non lo offrono. Utilizziamo l'infrastruttura cloud di Mail.ru Cloud Solutions e abbiamo scritto un autoscalatore utilizzando l'API MCS. Poiché insegniamo a lavorare con i dati, abbiamo deciso di mostrare come è possibile scrivere un simile autoscalatore per i propri scopi e utilizzarlo con il proprio cloud

Requisiti preliminari

Innanzitutto, è necessario avere un cluster Hadoop. Noi, per esempio, utilizziamo la distribuzione HDP.

Affinché i nodi possano essere aggiunti e rimossi rapidamente, è necessario avere una certa distribuzione dei ruoli tra i nodi.

  1. Nodo master. Non c'è molto da spiegare: il nodo principale del cluster, dove viene eseguito, ad esempio, il driver di Spark, se si utilizza la modalità interattiva.
  2. Data nodo. Questo è il nodo in cui sono memorizzati i dati su HDFS e dove avvengono anche i calcoli.
  3. Nodo di calcolo. Questo è il nodo in cui non viene memorizzato nulla su HDFS, ma dove avvengono i calcoli.

Un punto importante. L'autoscaling avverrà grazie ai nodi di terzo tipo. Se iniziate a rimuovere e aggiungere nodi di secondo tipo, la velocità di risposta sarà notevolmente ridotta: la decommissioning e il recommissioning richiederanno ore nel vostro cluster. Questo, ovviamente, non è ciò che ci si aspetta dall'autoscaling. Quindi non tocchiamo i nodi di primo e secondo tipo. Questi rappresenteranno un cluster minimamente vitale che esisterà per l'intera durata del programma.

Quindi, il nostro autoscaler è scritto in Python 3, utilizza l'API di Ambari per gestire i servizi del cluster, usa l'API di Mail.ru Cloud Solutions (MCS) per avviare e fermare le macchine.

Architettura della soluzione

  1. Modulo autoscaler.py. Contiene tre classi: 1) funzioni per lavorare con Ambari, 2) funzioni per lavorare con MCS, 3) funzioni collegate direttamente alla logica di funzionamento dell'autoscaler.
  2. Script observer.py. Essenzialmente consiste in diverse regole: quando e in quali momenti chiamare le funzioni dell'autoscaler.
  3. Il file con i parametri di configurazione config.py. Contiene, per esempio, un elenco dei nodi consentiti per l'autoscaling e altri parametri che influenzano, ad esempio, quanto tempo aspettare dal momento in cui viene aggiunto un nuovo nodo. Inoltre, ci sono anche i timestamp di inizio delle lezioni, affinché la configurazione massima consentita del cluster sia avviata prima della lezione.

Ora diamo un'occhiata ai frammenti di codice presenti nei primi due file.

1. Modulo autoscaler.py

Classe Ambari

Questo è un piccolo frammento di codice che contiene la classe 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

In alto, per esempio, puoi vedere l'implementazione della funzione stop_all_services, che arresta tutti i servizi sul nodo desiderato del cluster.

All'ingresso della classe Ambari si passano:

  • ambari_url, ad esempio del tipo 'http://localhost:8080/api/v1/clusters/',
  • cluster_name è il nome del tuo cluster in Ambari,
  • headers = {'X-Requested-By': 'ambari'}
  • e all'interno auth ci sono il tuo nome utente e la tua password per Ambari: auth = ('login', 'password').

La funzione stessa consiste in non più di un paio di chiamate attraverso l'API REST di Ambari. Dal punto di vista logico, inizialmente otteniamo l'elenco dei servizi attivi sul nodo, e poi chiediamo a questo cluster, su questo nodo, di portare i servizi dell'elenco nello stato INSTALLED. Le funzioni per avviare tutti i servizi, per mettere i nodi in stato Maintenance e altro, sono simili – si tratta semplicemente di alcune richieste tramite API.

Classe Mcs

Questo è un piccolo frammento di codice che contiene la classe 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

All'ingresso della classe Mcs passiamo l'id del progetto all'interno del cloud e l'id dell'utente, nonché la sua password. Nella funzione vm_turn_on Vogliamo attivare una delle macchine. La logica qui è un po' più complessa. All'inizio del codice, viene chiamata una funzione di tre altri: 1) dobbiamo ottenere un token, 2) dobbiamo convertire l'hostname nel nome della macchina in MCS, 3) otteniamo l'id di questa macchina. Dopo facciamo semplicemente una richiesta post e attiviamo questa macchina.

Ecco come si presenta la funzione per ottenere il token:

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

Classe Autoscaler

In questa classe sono presenti le funzioni relative alla logica di funzionamento.

Ecco un pezzo di codice di questa classe:

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('Richiesta di decommissione accettata: {0}'.format(flag1))
                    break
            while True:
                time.sleep(5)
                status3 = self.ambari.check_service(hostname, 'NODEMANAGER')
                if status3 == 'INSTALLED':
                    flag3 = True
                    logging.info('Nodemanager decommissionato: {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('Richiesta di manutenzione accettata: {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('La manutenzione è attiva: {0}'.format(flag4))
                    logging.info('Arresto servizi')
                    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('La VM è stata spenta: {0}'.format(flag5))
                    break
            if flag1 and flag2 and flag3 and flag4 and flag5:
                message = 'Successo'
                logging.info('Operazione di scale-down completata')
                logging.info('Periodo di raffreddamento iniziato. Attendere alcuni minuti')
        return message

In ingresso accettiamo classi Ambari e Mcs, una lista di nodi autorizzati per il scaling, nonché parametri di configurazione dei nodi: memoria e CPU allocate a ciascun nodo in YARN. Ci sono anche 2 parametri interni q_ram, q_cpu, che sono code. Attraverso di essi memorizziamo i valori dell'attuale carico del cluster. Se vediamo che negli ultimi 5 minuti c'è stato un carico elevato costante, decidiamo di aggiungere +1 nodo al cluster. Lo stesso vale per lo stato di sottocarico del cluster.

Nel codice sopra è mostrato un esempio di funzione che rimuove una macchina dal cluster e la arresta nel cloud. Inizialmente viene eseguito il decommissioning YARN Nodemanager, poi viene attivata la modalità Maintenance, poi fermiamo tutti i servizi sulla macchina e spegniamo la macchina virtuale nel cloud.

2. Script observer.py

Esempio di codice preso da lì:

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} è stato scalato con successo".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)

Qui controlliamo se le condizioni per l'aumento delle capacità del cluster sono soddisfatte e se ci sono macchine disponibili, otteniamo il nome di una di esse, la aggiungiamo al cluster e pubblichiamo un messaggio al riguardo nel Slack del nostro team. Successivamente, inizia cooldown_period, durante il quale non aggiungiamo né rimuoviamo nulla dal cluster, ma monitoriamo semplicemente il carico. Se si stabilizza e rimane all'interno del range ottimale di valori di carico, continuiamo semplicemente il monitoraggio. Se una node non è sufficiente, ne aggiungiamo un'altra.

In situazioni in cui abbiamo un'attività imminente, sappiamo già con certezza che una sola node non basterà, quindi avviamo immediatamente tutte le node disponibili e le manteniamo attive fino alla fine dell'attività. Questo avviene tramite un elenco di timestamp delle lezioni.

Conclusione

L'autoscaler è una soluzione valida e comoda per i casi in cui si osserva un carico irregolare nel cluster. Riuscite a ottenere contemporaneamente la configurazione desiderata del cluster per i carichi di picco, senza mantenere questo cluster durante i periodi di sotto-utilizzo, risparmiando così risorse. Inoltre, tutto ciò avviene in modo automatizzato senza il vostro intervento. L'autoscaler stesso non è altro che un insieme di richieste all'API del gestore del cluster e all'API del fornitore cloud, scritti secondo una logica specifica. Una cosa da tenere a mente è la suddivisione delle node in 3 tipi, come scritto in precedenza. E sarete felici.

Fonte: habr.com

Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server 🔥 Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server | ProHoster