Cómo crear tu propio escalador automático para un clúster

¡Hola! Enseñamos a las personas a trabajar con grandes volúmenes de datos. No se puede imaginar un programa educativo sobre grandes datos sin su propio clúster, donde todos los participantes trabajan juntos. Por esta razón, siempre lo tenemos en nuestro programa 🙂 Nos encargamos de su configuración, ajuste y administración, y los chicos ejecutan trabajos MapReduce y utilizan Spark allí.

En esta publicación, contaremos cómo resolvimos el problema de la carga desigual del clúster, escribiendo nuestro propio escalador automático utilizando la nube. Soluciones en la Nube de Mail.ru.

Problema

Nuestro clúster no se utiliza de manera completamente típica. La utilización es bastante desigual. Por ejemplo, hay clases prácticas cuando las 30 personas y el instructor acceden al clúster y comienzan a usarlo. O, nuevamente, hay días antes de la fecha límite cuando la carga aumenta considerablemente. En todo lo demás, el clúster opera en modo de subutilización.

Solución nº 1: mantener un clúster que pueda soportar las cargas máximas, pero que esté inactivo el resto del tiempo.

Solución nº 2: mantener un clúster pequeño, al que se le añadan nodos manualmente antes de las clases y durante los picos de carga.

Solución nº 3: mantener un clúster pequeño y escribir un escalador automático que vigile la carga actual del clúster y, utilizando diferentes API, añada y elimine nodos del clúster.

En esta publicación, hablaremos sobre la solución nº 3. Este tipo de escalador automático depende mucho de factores externos más que internos, y los proveedores a menudo no lo ofrecen. Utilizamos la infraestructura en la nube de Mail.ru Cloud Solutions y escribimos un escalador automático utilizando la API de MCS. Dado que enseñamos a trabajar con datos, decidimos mostrar cómo puedes escribir un escalador automático similar para tus objetivos y usarlo con tu propia nube.

Requisitos previos

En primer lugar, debes tener un clúster de Hadoop. Nosotros, por ejemplo, utilizamos la distribución HDP.

Para que los nodos puedan añadirse y eliminarse rápidamente, debes tener una distribución de roles específica en los nodos.

  1. Nodo maestro. Aquí no hace falta explicar demasiado: el nodo principal del clúster, donde se ejecuta, por ejemplo, el controlador de Spark, si usas el modo interactivo.
  2. Nodo de datos. Este es el nodo en el que almacenas los datos en HDFS y donde se realizan los cálculos.
  3. Nodo de cálculo. Esta es una nodo en la que no se almacena nada en HDFS, pero se realizan cálculos.

Un punto importante. El escalado automático se llevará a cabo a expensas de nodos de tercer tipo. Si comienzas a quitar y agregar nodos de segundo tipo, la velocidad de respuesta será muy baja: la descomisión y la recomisión tardarán horas en tu clúster. Esto, por supuesto, no es lo que esperas del escalado automático. Es decir, no tocamos los nodos de primer y segundo tipo. Representarán un clúster mínimamente viable, que existirá durante toda la duración del programa.

Así que nuestro escalador automático está escrito en Python 3, utiliza la API de Ambari para gestionar los servicios del clúster, utiliza la API de Mail.ru Cloud Solutions (MCS) para iniciar y detener máquinas.

Arquitectura de la solución

  1. Módulo autoscaler.py. Contiene tres clases: 1) funciones para trabajar con Ambari, 2) funciones para trabajar con MCS, 3) funciones relacionadas directamente con la lógica de funcionamiento del escalador automático.
  2. Script observer.py. Consiste esencialmente en diferentes reglas: cuándo y en qué momentos llamar a las funciones del escalador automático.
  3. Archivo con los parámetros de configuración config.py. Allí se encuentra, por ejemplo, la lista de nodos permitidos para el escalado automático y otros parámetros que influyen, por ejemplo, en cuánto tiempo esperar desde que se agregó un nuevo nodo. También contiene las marcas de tiempo de inicio de las clases, para que antes de la clase esté configurada la máxima configuración permitida del clúster.

Ahora veamos fragmentos de código que se encuentran dentro de los primeros dos archivos.

1. Módulo autoscaler.py

Clase Ambari

Así se ve un fragmento de código que contiene la clase 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":"Detener todos los componentes de host",
                "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 = 'Solicitud aceptada'
        else:
            message = req.status_code
        return message

Arriba, como ejemplo, se puede ver la implementación de la función stop_all_services, que detiene todos los servicios en el nodo requerido del clúster.

Al constructor de la clase Ambari se le pasan:

  • ambari_url, por ejemplo, del tipo 'http://localhost:8080/api/v1/clusters/',
  • cluster_name – el nombre de su clúster en Ambari,
  • headers = {'X-Requested-By': 'ambari'}
  • y dentro de auth está su nombre de usuario y contraseña de Ambari: auth = ('nombre_de_usuario', 'contraseña').

La función en sí consiste en no más que un par de llamadas a través de la API REST a Ambari. Desde el punto de vista lógico, primero obtenemos una lista de los servicios en ejecución en el nodo, y luego pedimos que en este clúster y en este nodo, los servicios de la lista se cambien al estado INSTALLED. Las funciones para iniciar todos los servicios, para cambiar los nodos al estado Mantenimiento y otras funcionan de manera similar: son simplemente varias solicitudes a través de la API.

Clase Mcs

Así se ve un fragmento de código que contiene la clase 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

Al constructor de la clase Mcs pasamos el id del proyecto dentro de la nube y el id del usuario, así como su contraseña. En la función vm_turn_on Queremos encender una de las máquinas. La lógica aquí es un poco más complicada. Al inicio del código se llaman a tres funciones diferentes: 1) necesitamos obtener un token, 2) necesitamos convertir el hostname en el nombre de la máquina en MCS, 3) obtener el id de esta máquina. A continuación, simplemente hacemos una solicitud POST y lanzamos esta máquina.

Así es como se ve la función para obtener el 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

Clase Autoscaler

En esta clase se encuentran las funciones relacionadas con la lógica de funcionamiento.

Así es como se ve un fragmento del código de esta clase:

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('Solicitud de descomisión aceptada: {0}'.format(flag1))
                    break
            while True:
                time.sleep(5)
                status3 = self.ambari.check_service(hostname, 'NODEMANAGER')
                if status3 == 'INSTALLED':
                    flag3 = True
                    logging.info('Nodemaneger descomisionado: {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('Solicitud de mantenimiento aceptada: {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('El mantenimiento está activado: {0}'.format(flag4))
                    logging.info('Deteniendo servicios')
                    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á apagada: {0}'.format(flag5))
                    break
            if flag1 and flag2 and flag3 and flag4 and flag5:
                message = 'Éxito'
                logging.info('Escalado hacia abajo completado')
                logging.info('Se ha iniciado el período de enfriamiento. Espere varios minutos')
        return message

Recibimos clases Ambari y Mcs, una lista de nodos permitidos para el escalado, así como parámetros de configuración de nodos: memoria y CPU asignadas al nodo en YARN. También hay 2 parámetros internos q_ram, q_cpu, que son colas. Con ellos almacenamos los valores de la carga actual del clúster. Si vemos que durante los últimos 5 minutos ha habido una carga alta constante, tomamos la decisión de agregar +1 nodo al clúster. Lo mismo es cierto para el estado de subcarga del clúster.

El código anterior muestra un ejemplo de una función que elimina una máquina del clúster y la detiene en la nube. Primero, se realiza la descomisión YARN Nodemanager, luego se activa el modo Mantenimiento, luego detenemos todos los servicios en la máquina y apagamos la máquina virtual en la nube.

2. El script observer.py

Ejemplo de código desde allí:

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} ha sido escalado exitosamente".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)

En él verificamos si se cumplen las condiciones para aumentar la capacidad del clúster y si hay máquinas en reserva, obtenemos el nombre de host de una de ellas, la añadimos al clúster y publicamos un mensaje sobre esto en Slack de nuestro equipo. Luego se inicia cooldown_period, cuando no añadimos ni quitamos nada del clúster, simplemente monitorizamos la carga. Si se estabiliza y se encuentra dentro del rango de valores óptimos de carga, continuamos la monitorización. Si una sola nodo no fue suficiente, añadimos otra.

En los casos en que tenemos una ocupación próxima, ya sabemos con certeza que una sola nodo no será suficiente, por lo que iniciamos de inmediato todas las nodos libres y las mantenemos activas hasta el final de la ocupación. Esto se realiza mediante una lista de marcas de tiempo de ocupaciones.

Conclusión

El autoescalador es una solución útil y conveniente para aquellos casos en los que hay una carga desigual en el clúster. Al mismo tiempo, logras la configuración adecuada del clúster para cargas picos y no mantienes este clúster durante los períodos de baja carga, ahorrando costos. Además, todo esto ocurre de manera automatizada sin tu intervención. El autoescalador en sí no es más que un conjunto de solicitudes a la API del gestor del clúster y a la API del proveedor de la nube, escritas según una lógica determinada. Lo que debes recordar con certeza es la división de nodos en 3 tipos, como hemos mencionado anteriormente. Y serás feliz.

Fuente: habr.com

Compra un hosting fiable para sitios web con protección contra DDoS, servidores VPS VDS 🔥 Compra un hosting fiable para sitios web con protección contra DDoS, servidores VPS VDS | ProHoster