Airflow — инструмент за удобно и бързо разработване и поддържане на batch-процеси за обработка на данни

Airflow — инструмент за удобно и бързо разработване и поддържане на batch-процеси за обработка на данни

Здравей, Хабър! В тази статия искам да ти разкажа за един чудесен инструмент за разработка на batch-процеси за обработка на данни, например, в инфраструктурата на корпоративен DWH или твоето DataLake. Ще става дума за Apache Airflow (по-нататък Airflow). Той е несправедливо пренебрегнат на Хабър и в основната част ще се опитам да те убедя, че поне Airflow си струва да разгледаш при избора на планировчик за твоите ETL/ELT процеси.

Предишния път написах серия статии на тема DWH, когато работех в Тинькофф Банка. Сега съм част от екипа на Mail.Ru Group и се занимавам с развитието на платформата за анализ на данни в игровия сектор. Всъщност, с течение на времето, ще споделяме тук с екипа новини и интересни решения относно нашата платформа за аналитика на данни.

Пролог

И така, да започнем. Какво е Airflow? Това е библиотека (или набор от библиотеки) за разработка, планиране и мониторинг на работни процеси. Основната особеност на Airflow е: за описанието (разработката) на процесите се използва код на Python. Оттук идват многобройните предимства за организиране на твоя проект и разработка: по същество, твоя (например) ETL проект — е просто Python проект и можеш да го организираш както желаеш, в зависимост от особеностите на инфраструктурата, размера на екипа и други изисквания. Инструментално всичко е просто. Използвай, например, PyCharm + Git. Това е прекрасно и много удобно!

Сега нека разгледаме основните същности на Airflow. Като осъзнаеш тяхната същност и предназначение, ще организираш архитектурата на процесите оптимално. Вероятно основната същност е Directed Acyclic Graph (по-нататък DAG).

DAG

DAG е някакво смислово обединение на твоите задачи, които искаш да изпълниш в строго определена последователност по определен график. Airflow предлага удобен web интерфейс за работа с DAG-ове и други същности:

Airflow — инструмент за удобно и бързо разработване и поддържане на batch-процеси за обработка на данни

DAG може да изглежда по следния начин:

Airflow — инструмент за удобно и бързо разработване и поддържане на batch-процеси за обработка на данни

Разработчикът, проектирайки DAG, задава набор от оператори, на които ще се основават задачите вътре в DAG-а. Тук стигаме до още една важна същност: Airflow Operator.

Оператори

Операторът е същност, на базата на която се създават екземпляри на задания, където е описано какво ще се случва по време на изпълнението на екземпляра на заданието. Релизите на Airflow от GitHub вече съдържат набор от оператори, готови за употреба. Примери:

  • BashOperator — оператор за изпълнение на bash команда.
  • PythonOperator — оператор за извикване на Python код.
  • EmailOperator — оператор за изпращане на имейл.
  • HTTPOperator — оператор за работа с http заявки.
  • SqlOperator — оператор за изпълнение на SQL код.
  • Sensor — оператор, който изчаква събитие (появата на определено време, необходим файл, ред в базата данни, отговор от API и т.н.).

Има и по-специфични оператори: DockerOperator, HiveOperator, S3FileTransferOperator, PrestoToMysqlOperator, SlackOperator.

Можете също така да разработвате оператори в зависимост от вашите нужди и да ги използвате в проекта. Например, създадохме MongoDBToHiveViaHdfsTransfer, оператор за експортиране на документи от MongoDB в Hive, и няколко оператора за работа с ClickHouse: CHLoadFromHiveOperator и CHTableLoaderOperator. Всъщност, когато в проекта възникне често използван код, базиран на основни оператори, можете да се замислите дали да не го обедините в нов оператор. Това ще улесни последващата разработка и ще обогати библиотеката ви от оператори в проекта.

Следва да се изпълняват всички тези экземпляри на задачи, и сега ще говорим за планировчика.

Два забележителни изменения в планирането (и двете в алфа версия):

Планировчикът на задачи в Airflow е изграден на Celery. Celery е Python библиотека, която позволява организиране на опашка, плюс асинхронно и разпределено изпълнение на задачи. От страна на Airflow всички задачи се разпределят в пули. Пулите се създават ръчно. Обикновено те имат за цел да ограничат натоварването при работа с източника или да типизират задачите в DWH. Пулите могат да се управляват чрез уеб интерфейса:

Airflow — инструмент за удобно и бързо разработване и поддържане на batch-процеси за обработка на данни

Всеки пул има ограничение по броя на слотовете. При създаване на DAG, му се задава пул:

ALERT_MAILS =  Variable.get("gv_mail_admin_dwh")
DAG_NAME = 'dma_load'
OWNER = 'Vasya Pupkin'
DEPENDS_ON_PAST = True
EMAIL_ON_FAILURE = True
EMAIL_ON_RETRY = True
RETRIES = int(Variable.get('gv_dag_retries'))
POOL = 'dma_pool'
PRIORITY_WEIGHT = 10

start_dt = datetime.today() - timedelta(1)
start_dt = datetime(start_dt.year, start_dt.month, start_dt.day)

default_args = {
    'owner': OWNER,
    'depends_on_past': DEPENDS_ON_PAST,
    'start_date': start_dt,
    'email': ALERT_MAILS,
    'email_on_failure': EMAIL_ON_FAILURE,
    'email_on_retry': EMAIL_ON_RETRY,
    'retries': RETRIES,
    'pool': POOL,
    'priority_weight': PRIORITY_WEIGHT
}
dag = DAG(DAG_NAME, default_args=default_args)
dag.doc_md = __doc__

Пулът, зададен на ниво DAG, може да бъде заменен на ниво задача.
За планирането на всички задачи в Airflow отговаря отделен процес — Scheduler. Всъщност, Scheduler се занимава с механиката на поставянето на задачите за изпълнение. Задачата, преди да бъде изпълнена, преминава през няколко етапа:

  1. В DAG'а са изпълнени предишните задачи, новата може да се постави в опашката.
  2. Опашката се сортира в зависимост от приоритета на задачите (приоритетите също могат да се управляват) и, ако в пула има свободен слот, задачата може да бъде взета за работа.
  3. Ако има свободен worker celery, задачата се насочва към него; започва работа, която сте програмирали в задачата, използвайки този или онзи оператор.

Достатъчно просто.

Scheduler работи върху множество DAG'ове и всички задачи вътре в DAG'овете.

За да започне Scheduler работа с DAG'а, на DAG'а трябва да се зададе разписание:

dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='@hourly')

Има набор от готови preset'ове: @once, @hourly, @daily, @weekly, @monthly, @yearly.

Също така можете да използвате cron-изрази:

dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='*\/10 * * * *')

Execution Date

За да разберете как работи Airflow, важно е да разберете какво е Execution Date за DAG'а. В Airflow, DAG има измерение Execution Date, т.е. в зависимост от разписанието на работа на DAG'а, се създават екземпляри на задачите за всяка Execution Date. И за всяка Execution Date задачите могат да се изпълняват повторно — или, например, DAG'а може да работи едновременно в няколко Execution Date. Това е илюстрирано тук:

Airflow — инструмент за удобно и бързо разработване и поддържане на batch-процеси за обработка на данни

За съжаление (а може би и за щастие: зависи от ситуацията), ако се поправи реализацията на задачата в DAG'а, то изпълнението в предишните Execution Date ще се проведе с оглед на корекциите. Това е добре, ако трябва да се преработят данните в минали периоди с нов алгоритъм, но и лошо, защото се губи възпроизводимостта на резултата (разбира се, никой не пречи да се върне нужната версия на изходния код от Git и да се изчисли всичко, както трябва, веднъж).

Генерация на задачи

Реализацията на DAG'а е код на Python, така че имаме много удобен начин да съкратим обема код при работа, например, с шардировни източници. Да приемем, че имате три шарда MySQL като източник, трябва да влезете в всеки от тях и да извлечете някакви данни. При това независимо и паралелно. Кодът на Python в DAG'а може да изглежда така:

connection_list = lv.get('connection_list')

export_profiles_sql = '''
SELECT
  id,
  user_id,
  nickname,
  gender,
  {{params.shard_id}} as shard_id
FROM profiles
'''

for conn_id in connection_list:
    export_profiles = SqlToHiveViaHdfsTransfer(
        task_id='export_profiles_from_' + conn_id,
        sql=export_profiles_sql,
        hive_table='stg.profiles',
        overwrite=False,
        tmpdir='\/data\/tmp',
        conn_id=conn_id,
        params={'shard_id': conn_id[-1:], },
        compress=None,
        dag=dag
    )
    export_profiles.set_upstream(exec_truncate_stg)
    export_profiles.set_downstream(load_profiles)

DAG'а изглежда така:

Airflow — инструмент за удобно и бързо разработване и поддържане на batch-процеси за обработка на данни

В този случай можете да добавите или премахнете шард, просто като промените настройките и обновите DAG. Удобно!

Можете да използвате и по-сложна генерация на код, например да работите с източници под формата на БД или да опишете таблична структура, алгоритъм за работа с таблицата, и с оглед особеностите на инфраструктурата DWH да генерирате процес за зареждане на N таблици в вашето хранилище. Или, например, работа с API, което не поддържа работа с параметър под формата на списък, можете да генерирате N задачи в DAG-а по този списък, да ограничите паралелността на заявките в API с пул и да изтеглите необходимите данни от API. Гъвкаво!

Репозиторий

В Airflow има собствен бекенд репозиторий, БД (може да бъде MySQL или Postgres, при нас е Postgres), в който се съхраняват състояния на задачи, DAG-ове, настройки на връзки, глобални променливи и т.н. Тук бих искал да кажа, че репозиторият в Airflow е много прост (около 20 таблици) и удобен, ако искате да изградите собствен процес. Спомням си 100500 таблици в репозиторий на Informatica, които трябваше да разглеждам дълго време, преди да разбера как да изградя запитване.

Мониторинг

Като се има предвид простотата на репозитория, можете сами да изградите удобен за вас процес за мониторинг на задачите. Ние използваме бележник в Zeppelin, където следим състоянието на задачите:

Airflow — инструмент за удобно и бързо разработване и поддържане на batch-процеси за обработка на данни

Това може да бъде и уеб интерфейсът на самия Airflow:

Airflow — инструмент за удобно и бързо разработване и поддържане на batch-процеси за обработка на данни

Кодът на Airflow е отворен, затова добавихме аларми в Telegram. Всеки работещ инстанс на задача, ако възникне грешка, изпраща спам в групата в Telegram, където е целият екип по разработка и поддръжка.

Получаваме чрез Telegram бързи реакции (ако е необходимо), а чрез Zeppelin - обща картина на задачите в Airflow.

Итого

Airflow е преди всичко open source, и не трябва да очаквате чудеса от него. Бъдете готови да инвестирате време и усилия, за да изградите работещо решение. Целта е постижима, повярвайте, заслужава си. Скоростта на разработка, гъвкавостта, простотата на добавяне на нови процеси - ще ви хареса. Разбира се, трябва да се обърне много внимание на организацията на проекта, стабилността на самия Airflow: чудеса няма.

В момента нашият Airflow ежедневно обработва около 6500 задачи. По характеру те са достатъчно различни. Има задачи за зареждане на данни в основното DWH от множество различни и много специфични източници, има задачи за изчисление на витрини вътре в основното DWH, има задачи за публикуване на данни в бързо DWH, има много-много различни задачи — и Airflow всичките ги обработва ден след ден. Ако говорим с цифри, то това 2,3 хиляди ELT задачи с различна сложност вътре в DWH (Hadoop), около 2,5 стотин бази данни източници, това е екип от 4-ма ETL разработчици, които се разделят на ETL процесинг на данни в DWH и на ELT процесинг на данни вътре в DWH и разбира се още един администратор, който се занимава с инфраструктурата на услугата.

Планове за бъдещето

Броят на процесите неизбежно расте, а основното, с което ще се занимаваме в частта за инфраструктура на Airflow, е мащабиране. Искаме да изградим кластер Airflow, да отделим няколко ядра за worker-ите на Celery и да направим дублираща се глава с процеси за планиране на задачи и хранилище.

Епилог

Това, разбира се, не е всичко, което бих искал да разкажа за Airflow, но основните моменти се опитах да осветя. Апетитът идва с яденето, опитайте — и ще ви хареса 🙂

Източник: habr.com

Купете надежден хостинг за сайтове с защита от DDoS, VPS VDS сървъри 🔥 Купете надежден хостинг за сайтове с защита от DDoS, VPS VDS сървъри | ProHoster