Comment créer votre propre auto-scaler pour un cluster

Bonjour ! Nous formons des personnes à travailler avec de grandes données. Il est impensable d'avoir un programme éducatif sur les grandes données sans son propre cluster où tous les participants peuvent collaborer. Pour cette raison, il est toujours présent dans notre programme 🙂 Nous nous occupons de sa configuration, de son optimisation et de son administration, tandis que les participants y lancent des travaux MapReduce et utilisent Spark.

Dans cet article, nous allons expliquer comment nous avons résolu le problème de la charge inégale du cluster en écrivant notre propre auto-scaleur, en utilisant le cloud. Mail.ru Cloud Solutions.

Le problème

Notre cluster n'est pas utilisé dans un mode typique. L'utilisation est fortement inégale. Par exemple, lors des cours pratiques, les 30 participants et l'instructeur se connectent au cluster et commencent à l'utiliser. Il y a aussi des jours avant les délais où la charge augmente considérablement. À tous les autres moments, le cluster fonctionne en mode sous-utilisé.

Solution n°1 – maintenir un cluster capable de supporter des charges de pointe, mais qui sera inactif le reste du temps.

Solution n°2 – maintenir un petit cluster dans lequel ajouter manuellement des nœuds avant les cours et durant les périodes de forte charge.

Solution n°3 – maintenir un petit cluster et écrire un auto-scaleur qui surveillera la charge actuelle du cluster et ajoutera ou supprimera des nœuds de manière autonome en utilisant différentes API.

Dans cet article, nous allons parler de la solution n°3. Cet auto-scaleur dépend fortement de facteurs externes, plutôt que internes, et les fournisseurs ne le proposent souvent pas. Nous utilisons l'infrastructure cloud Mail.ru Cloud Solutions et avons écrit un auto-scaleur en utilisant l'API MCS. Comme nous formons à la manipulation de données, nous avons décidé de montrer comment vous pouvez écrire un auto-scaleur similaire pour vos propres objectifs et l'utiliser avec votre cloud.

Prérequis

Tout d'abord, vous devez avoir un cluster Hadoop. Par exemple, nous utilisons la distribution HDP.

Pour que vos nœuds puissent être ajoutés et supprimés rapidement, vous devez avoir une certaine répartition des rôles parmi les nœuds.

  1. Nœud maître. Il n'est pas nécessaire de donner d'explications ici : c'est le nœud principal du cluster, sur lequel est lancé, par exemple, le driver de Spark, si vous utilisez le mode interactif.
  2. Nœud de données. C'est le nœud sur lequel vos données sont stockées sur HDFS et où les calculs ont également lieu.
  3. Nœud de calcul. C'est un nœud sur lequel rien n'est stocké sur HDFS, mais des calculs y sont effectués.

Point important. L'auto-scaling se fera grâce aux nœuds de type trois. Si vous commencez à retirer et ajouter des nœuds de type deux, la réactivité sera considérablement réduite - le décommissionnement et le recommissionnement prendra des heures dans votre cluster. Ce n'est évidemment pas ce que l'on attend d'un auto-scaling. Donc, nous ne touchons pas aux nœuds de type un et deux. Ils représenteront un cluster minimement viable qui existera tout au long de la durée du programme.

Ainsi, notre auto-scaler est écrit en Python 3, utilise l'API d'Ambari pour gérer les services du cluster, utilise l'API de Mail.ru Cloud Solutions (MCS) pour démarrer et arrêter les machines.

Architecture de la solution

  1. Module autoscaler.py. Il contient trois classes : 1) des fonctions pour travailler avec Ambari, 2) des fonctions pour travailler avec MCS, 3) des fonctions liées directement à la logique de fonctionnement de l'auto-scaler.
  2. Script observer.py. En substance, il se compose de différentes règles : quand et à quels moments appeler les fonctions de l'auto-scaler.
  3. Le fichier avec les paramètres de configuration config.py. Il contient, par exemple, la liste des nœuds autorisés pour l'auto-scaling et d'autres paramètres influençant, par exemple, le temps à attendre depuis l'ajout d'un nouveau nœud. Il contient également les timestamps de début des cours, afin que la configuration maximale autorisée du cluster soit lancée avant le cours.

Examinons maintenant des morceaux de code présents dans les deux premiers fichiers.

1. Module autoscaler.py

Classe Ambari

Voici un extrait de code contenant 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":"Arrêter tous les composants d'hôte",
                "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 = 'Demande acceptée'
        else:
            message = req.status_code
        return message

Vous pouvez voir ci-dessus l'implémentation de la fonction stop_all_services, qui arrête tous les services sur le nœud souhaité du cluster.

En entrant dans la classe Ambari vous fournissez :

  • ambari_url, par exemple sous la forme 'http://localhost:8080/api/v1/clusters/',
  • cluster_name – le nom de votre cluster dans Ambari,
  • headers = {'X-Requested-By': 'ambari'}
  • et à l'intérieur auth se trouve votre identifiant et votre mot de passe pour Ambari : auth = ('login', 'password').

La fonction elle-même consiste simplement en quelques appels via l'API REST d'Ambari. Logiquement, nous commençons par obtenir la liste des services en cours d'exécution sur le nœud, puis nous demandons sur le cluster donné, sur le nœud donné, de changer l'état des services de la liste en INSTALLED. Les fonctions pour démarrer tous les services, pour mettre les nœuds en état Maintenance et d'autres ressemblent à cela – ce sont juste plusieurs demandes via l'API.

Classe Mcs

Voici un extrait de code contenant 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

En entrant dans la classe Mcs nous transmettons l'identifiant du projet dans le cloud ainsi que l'identifiant de l'utilisateur et son mot de passe. Dans la fonction vm_turn_on Nous souhaitons activer l'une des machines. La logique ici est un peu plus complexe. Au début du code, trois autres fonctions sont appelées : 1) nous devons obtenir un jeton, 2) nous devons convertir le nom d'hôte en nom de la machine dans MCS, 3) obtenir l'ID de cette machine. Ensuite, nous effectuons simplement une requête POST pour lancer cette machine.

Voici la fonction qui permet d'obtenir le jeton :

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

Cette classe contient des fonctions liées à la logique de fonctionnement.

Voici un extrait de code de cette 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('Demande de décommission acceptée : {0}'.format(flag1))
                    break
            while True:
                time.sleep(5)
                status3 = self.ambari.check_service(hostname, 'NODEMANAGER')
                if status3 == 'INSTALLED':
                    flag3 = True
                    logging.info('Nodemaneger décommissionné : {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('Demande de maintenance acceptée : {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 maintenance est active : {0}'.format(flag4))
                    logging.info('Arrêt des 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('La VM est éteinte : {0}'.format(flag5))
                    break
            if flag1 and flag2 and flag3 and flag4 and flag5:
                message = 'Succès'
                logging.info('Scale-down terminé')
                logging.info('La période de refroidissement a commencé. Attendez quelques minutes')
        return message

En entrée, nous avons des classes Ambari et Mcs, une liste de nœuds autorisés pour le scaling, ainsi que les paramètres de configuration des nœuds : la mémoire et le CPU alloués à chaque nœud dans YARN. Il y a aussi 2 paramètres internes q_ram, q_cpu, qui sont des files d'attente. Grâce à eux, nous stockons les valeurs de charge actuelle du cluster. Si nous observons qu'il y a eu une charge élevée stable durant les 5 dernières minutes, nous décidons d'ajouter +1 nœud au cluster. Il en va de même pour les situations de sous-charge du cluster.

Dans le code ci-dessus, un exemple de fonction est donné, qui retire une machine du cluster et l'arrête dans le cloud. Au départ, nous effectuons le décommissionnement YARN Nodemanager, ensuite, nous activons le mode Maintenance, puis nous arrêtons tous les services sur la machine et éteignons la machine virtuelle dans le cloud.

2. Le script observer.py

Exemple de code provenant de 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} a été mis à l'échelle avec succès".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)

Ici, nous vérifions si les conditions sont réunies pour augmenter la capacité du cluster et s'il y a des machines disponibles, nous obtenons le nom d'hôte de l'une d'elles, l'ajoutons au cluster et publions un message à ce sujet dans Slack de notre équipe. Ensuite, un cooldown_period, où nous n'ajoutons ni ne retirons quoi que ce soit du cluster, mais surveillons simplement la charge. Si celle-ci s'est stabilisée et reste dans le couloir des valeurs de charge optimales, nous continuons simplement la surveillance. Si une seule nœud ne suffit pas, nous en ajoutons une autre.

Pour les cas où nous savons d'avance qu'une seule nœud ne suffira pas, nous lançons immédiatement tous les nœuds disponibles et les maintenons actifs jusqu'à la fin de la session. Cela se fait à l'aide d'une liste de timestamps de sessions.

Conclusion

L'autoscaler est une solution pratique et efficace pour les cas où vous observez une charge inégale sur le cluster. Vous obtenez simultanément la configuration souhaitée du cluster pour les charges maximales tout en évitant de maintenir ce cluster pendant les périodes de sous-charge, ce qui vous permet d'économiser des fonds. De plus, tout cela se fait automatiquement sans votre intervention. L'autoscaler lui-même n'est rien d'autre qu'un ensemble de requêtes à l'API du gestionnaire de cluster et à l'API du fournisseur de cloud, écrites selon une certaine logique. Ce dont il faut absolument se souvenir, c'est de la division des nœuds en 3 types, comme nous l'avons écrit précédemment. Et vous serez heureux.

Source : habr.com

Acheter un hébergement fiable pour les sites avec protection DDoS, serveurs VPS VDS 🔥 Acheter un hébergement fiable pour les sites avec protection DDoS, serveurs VPS VDS | ProHoster