Здравейте! Ние обучаваме хората как да работят с големи данни. Невъзможно е да си представим образователна програма за големи данни без собствен клъстер, на който всички участници да работят заедно. Поради тази причина в нашата програма той винаги е наличен 🙂 Ние се занимаваме с неговото настройване, оптимизиране и администриране, а участниците в него директно стартират MapReduce задачи и използват Spark.
В този пост ще разкажем как решихме проблема с неравномерното натоварване на клъстера, като написахме собствен автоскейлър, използвайки облака. .
Проблема
Нашият клъстер не се използва в типичен режим. Утилизацията е силно неравномерна. Например, има практически занятия, когато всички 30 души и преподавателят влизат в клъстера и започват да го използват. Или отново има дни преди крайните срокове, когато натоварването значително нараства. През останалото време клъстерът работи в режим на недоизползване.
Решение №1 – да държим клъстер, който ще може да издържи пиковите натоварвания, но ще остане бездействуващ през останалото време.
Решение №2 – да поддържаме малък клъстер, в който ръчно да добавяме възли преди занятията и по време на пиковите натоварвания.
Решение №3 – да поддържаме малък клъстер и да напишем автоскейлър, който да следи текущото натоварване на клъстера и самостоятелно, използвайки различни API, да добавя и премахва възли от клъстера.
В този пост ще говорим за решение №3. Такъв автоскейлър зависи в значителна степен от външни фактори, а не от вътрешни, и доставчиците често не го предоставят. Ние използваме облачната инфраструктура на Mail.ru Cloud Solutions и написахме автоскейлър, използвайки API MCS. Тъй като обучаваме работа с данни, решихме да покажем как можете да напишете подобен автоскейлър за вашите нужди и да го използвате с вашия облак.
Предварителни условия
На първо място, трябва да имате Hadoop клъстер. Ние, например, използваме дистрибуция на HDP.
За да можете бързо да добавяте и премахвате възли, трябва да имате определено разпределение на ролите между възлите.
- Майстор-възел. Тук не се нуждаете от много обяснения: главният възел на клъстера, на който стартира, например, драйверът на Spark, ако използвате интерактивен режим.
- Данни-възел. Това е възелът, на който съхранявате данни в HDFS и на него също се извършват изчисления.
- Изчислителна възел. Това е възел, на който не съхранявате нищо на HDFS, но на него се извършват изчисления.
Важно. Автоматичното мащабиране ще се извършва с помощта на третия тип възли. Ако започнете да изтегляте и добавяте възли от втория тип, реакционната скорост ще бъде значително ниска - декомишенът и рекомишенът ще отнемат часове в кластера ви. Това определено не е нещо, което очаквате от автоматичното мащабиране. Тоест, първият и вторият тип възли остават непроменени. Те ще представляват минимално жизнеспособен клъстер, който ще съществува през целия период на програмата.
И така, нашият автоматичен скелер е написан на Python 3, използва Ambari API за управление на услугите в клъстера, използва (MCS) за стартиране и спиране на машините.
Архитектура на решението
- Модулът
autoscaler.py. В него са описани три класа: 1) функции за работа с Ambari, 2) функции за работа с MCS, 3) функции, свързани пряко с логиката на работа на автоматичния скелер. - Скрипт
observer.py. По същество се състои от различни правила: кога и в какъв момент да се извикват функциите на автоматичния скелер. - Файл с конфигурационни параметри
config.py. Съдържа, например, списък на възлите, разрешени за автоматично мащабиране, и други параметри, влияещи, например, на това колко време да се изчака от момента, когато е добавен нов възел. Там също така са времевите печати за началото на занятията, за да гарантират, че преди занятията е активирана максимално разрешената конфигурация на клъстера.
Нека сега разгледаме фрагменти от кода, намиращи се в първите два файла.
1. Модул autoscaler.py
Клас Ambari
Така изглежда фрагмент от кода, съдържащ класа Ambari:
клас Ambari:
def __init__(self, ambari_url, cluster_name, заголовки, auth):
self.ambari_url = ambari_url
self.cluster_name = cluster_name
self.заголовки = заголовки
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, заголовки=self.заголовки, auth=self.auth)
services = req0.json()['host_components']
services_list = list(map(lambda x: x['HostRoles']['component_name'], services))
data = {
"RequestInfo": {
"context":"Остановить все компоненты хоста",
"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), заголовки=self.заголовки, auth=self.auth)
if req.status_code in [200, 201, 202]:
message = 'Запрос принят'
else:
message = req.status_code
return messageПо вышеуказанному примеру можно увидеть реализацию функции stop_all_services, которая останавливает все сервисы на нужном узле кластера.
В качестве аргументов классу Ambari вы передаете:
ambari_url, например, вида'http://localhost:8080/api/v1/clusters/',cluster_name– название вашего кластера в Ambari,headers = {'X-Requested-By': 'ambari'}- и внутри
authнаходятся ваш логин и пароль от Ambari:auth = ('логин', 'пароль').
Сама функция фактически представляет собой выполненные обращения через REST API к Ambari. С точки зрения логики, мы сначала получаем список запущенных сервисов на узле, а затем просим этот кластер, на данном узле перевести сервисы из списка в состояние INSTALLED. Функции, которые запускают все сервисы, переводят узлы в состояние Maintenance и другие выглядят аналогично – это просто несколько запросов через API.
Класс Mcs
Така изглежда фрагмент от кода, съдържащ класа 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'
заголовки = {
'X-Auth-Token': '{0}'.format(self.token),
'Content-Type': 'application/ json'
}
data = {'os-start' : 'null'}
mcs = requests.post(mcs_url1, data=json.dumps(data), заголовки=заголовки)
return mcs.status_codeВ качестве аргументов классу Mcs мы передаем id проекта в облаке и id пользователя, а также его пароль. В функции vm_turn_on Искаме да активираме една от машините. Логиката тук е малко по-сложна. В началото на кода се извикват три други функции: 1) трябва да получим токен, 2) трябва да конвертираме hostname в името на машината в MCS, 3) да получим id на тази машина. След това просто правим post-заявка и стартираме тази машина.
Така изглежда самата функция за получаване на токен:
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
В този клас се съдържат функции, свързани със самата логика на работа.
Така изглежда част от кода на този клас:
клас 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('Заявка на декомишен принята: {0}'.format(flag1))
break
while True:
time.sleep(5)
status3 = self.ambari.check_service(hostname, 'NODEMANAGER')
if status3 == 'INSTALLED':
flag3 = True
logging.info('Nodemaneger декомишен: {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('Заявка на обслуживание принята: {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('Обслуживание включено: {0}'.format(flag4))
logging.info('Остановка сервисов')
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('Виртуальная машина выключена: {0}'.format(flag5))
break
if flag1 and flag2 and flag3 and flag4 and flag5:
message = 'Успех'
logging.info('Процесс уменьшения завершен')
logging.info('Период ожидания начался. Подождите несколько минут')
return messageНа вход получаваме класове Ambari и Mcs, списък от възлови точки, разрешени за скалиране, и конфигурационни параметри на възлите: памет и CPU, зададени за възел в YARN. Освен това има 2 вътрешни параметъра q_ram, q_cpu, представляващи опашки. Чрез тях съхраняваме стойностите на текущото натоварване на клъстера. Ако видим, че през последните 5 минути натоварването е било стабилно повишено, решаваме, че трябва да добавим +1 възел в клъстера. Същото важи и за състоянието на недоразтоварване на клъстера.
В горния код е приведен пример за функция, която премахва машина от клъстера и я спира в облака. Първо се извършва декомишен YARN Nodemanager, след това се включва режим Maintenance, след това спираме всички услуги на машината и изключваме виртуалната машина в облака.
2. Скрипт observer.py
Примерен код оттам:
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} беше успешно мащабирано".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)В него проверяваме дали условията за увеличаване на мощностите на кластера са налице и дали имаме резервни машини, получаваме хостнейма на една от тях, добавяме я в клъстера и публикуваме съобщение за това в Slack на нашия екип. След това се стартира период на охлаждане, когато не добавяме и не премахваме нищо от клъстера, а просто наблюдаваме натоварването. Ако то се е стабилизирало и е в оптималните стойности на натоварване, просто продължаваме с наблюдението. Но ако една нода не е достатъчна, добавяме още една.
За случаите, когато имаме предстояща занимавка, вече знаем със сигурност, че една нода няма да е достатъчна, затова веднага стартираме всички свободни ноди и ги поддържаме активни до края на занимавката. Това се осъществява чрез списък с времеви отметки на заниманията.
Заключение
Автоскейлерът е добро и удобно решение за случаите, когато наблюдавате неравномерно натоварване на кластера. Вие едновременно постигате желаната конфигурация на кластера за пикова натовареност и в същото време не поддържате клъстера по време на недонасяне, спестявайки средства. Освен това, всичко това се случва автоматизирано без ваше участие. Самият автоскейлер не е нищо повече от набор от заявки към API на мениджъра на клъстера и API на облачния доставчик, записани в определена логика. За което е важно да помните – за разделянето на нодите на 3 типа, какво сме писали по-рано. И ще имате успех.
Източник: habr.com
