
Cześć, Habr! W tym artykule chcę opowiedzieć o jednym wspaniałym narzędziu do opracowywania procesów wsadowych przetwarzania danych, na przykład w infrastrukturze korporacyjnego DWH lub w twoim DataLake. Mowa będzie o Apache Airflow (dalej Airflow). Niestety, nie jest on odpowiednio doceniany na Habra, a w głównej części postaram się przekonać was, że przynajmniej warto zwrócić uwagę na Airflow przy wyborze planera dla waszych procesów ETL/ELT.
Wcześniej pisałem serię artykułów na temat DWH, kiedy pracowałem w Tinkoff Bank. Teraz jestem częścią zespołu Mail.Ru Group i zajmuję się rozwojem platformy do analizy danych w obszarze gier. W miarę pojawiania się newsów i interesujących rozwiązań, my i zespół będziemy tu opowiadać o naszej platformie do analityki danych.
Prolog
Rozpocznijmy zatem. Czym jest Airflow? To biblioteka (albo ) do opracowywania, planowania i monitorowania procesów roboczych. Podstawową cechą Airflow jest to, że do opisywania (oprowadzania) procesów stosuje się kod w języku Python. Z tego wynika mnóstwo korzyści dla organizacji projektu i rozwoju: w zasadzie, wasz (na przykład) projekt ETL to po prostu projekt w Pythonie, który możecie organizować w sposób, jaki wam odpowiada, biorąc pod uwagę cechy infrastruktury, wielkość zespołu i inne wymagania. Z narzędziami wszystko jest proste. Używajcie na przykład PyCharm + Git. To wspaniałe i bardzo wygodne!
Teraz przyjrzyjmy się głównym jednostkom Airflow. Zrozumienie ich istoty i przeznaczenia pozwoli optymalnie zorganizować architekturę procesów. Najważniejszą jednostką jest Directed Acyclic Graph (dalej DAG).
DAG
DAG to pewne logiczne połączenie zadań, które chcecie wykonać w ściśle określonej kolejności zgodnie z ustalonym harmonogramem. Airflow oferuje wygodny interfejs webowy do pracy z DAG-ami i innymi jednostkami:

DAG może wyglądać w ten sposób:

Programista, projektując DAG, określa zestaw operatorów, na których będą oparte zadania w obrębie DAG-a. Tutaj przychodzimy do jeszcze jednej ważnej jednostki: Airflow Operator.
Operatory
Operator to jednostka, na podstawie której tworzone są instancje zadań, w których opisuje się, co ma się dziać podczas wykonywania instancji zadania. już zawierają zestaw gotowych do użycia operatorów. Przykłady:
- BashOperator — operator do wykonywania polecenia bash.
- PythonOperator — operator do wywoływania kodu Python.
- EmailOperator — operator do wysyłania e-maila.
- HTTPOperator — operator do pracy z zapytaniami http.
- SqlOperator — operator do wykonywania kodu SQL.
- Sensor — operator oczekiwania na zdarzenie (np. osiągnięcie określonego czasu, pojawienie się wymaganego pliku, wiersza w bazie danych, odpowiedzi z API itp.).
Istnieją bardziej specyficzne operatory: DockerOperator, HiveOperator, S3FileTransferOperator, PrestoToMysqlOperator, SlackOperator.
Możesz także rozwijać operatory, dostosowując je do swoich potrzeb i używając ich w projekcie. Na przykład stworzyliśmy MongoDBToHiveViaHdfsTransfer, operatora eksportującego dokumenty z MongoDB do Hive, oraz kilka operatorów do pracy z CHLoadFromHiveOperator i CHTableLoaderOperator. W zasadzie, gdy w projekcie pojawia się często używany kod, oparty na podstawowych operatorach, warto pomyśleć o stworzeniu nowego operatora. Ułatwi to dalszy rozwój, a dodatkowo wzbogacisz swoją bibliotekę operatorów w projekcie. Następnie wszystkie te instancje zadań muszą zostać wykonane, a teraz skupimy się na harmonogramie.
Harmonogram zadań w Airflow oparty jest na
Harmonogram
Celery Każda pula ma ograniczenie co do liczby slotów. Przy tworzeniu DAG-a przypisuje się mu pulę:

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 = 10start_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__
Pulę przypisaną na poziomie DAG-a można nadpisać na poziomie zadania.Za planowanie wszystkich zadań w Airflow odpowiada osobny proces — Scheduler. Sam Scheduler zajmuje się całą mechaniką ustalania zadań do wykonania. Zadanie, zanim trafi do wykonania, przechodzi przez kilka etapów:
За планировку всех задач в Airflow отвечает отдельный процесс — Scheduler. Собственно, Scheduler занимается всей механикой постановки задачек на исполнение. Задача, прежде чем попасть на исполнение, проходит несколько этапов:
- W DAG-u wykonano poprzednie zadania, nowe można ustawić w kolejce.
- Kolejka jest sortowana w zależności od priorytetu zadań (priorytetami również można zarządzać), a jeśli w puli jest wolny slot, zadanie można podjąć do pracy.
- Jeśli jest wolny worker celery, zadanie jest mu przekazywane; rozpoczyna się praca, którą zaprogramowałeś w zadaniu, korzystając z określonego operatora.
Szczerze mówiąc, to dość proste.
Scheduler działa na wielu DAG-ach i wszystkich zadaniach wewnątrz DAG-ów.
Aby Scheduler rozpoczął pracę z DAG-iem, DAG musi mieć określony harmonogram:
dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='@hourly')Istnieje zestaw gotowych presetów: @once, @hourly, @daily, @weekly, @monthly, @yearly.
Można również używać wyrażeń cron:
dag = DAG(DAG_NAME, default_args=default_args, schedule_interval='*\/10 * * * *')Data wykonania
Aby zrozumieć, jak działa Airflow, ważne jest, aby wiedzieć, czym jest Data wykonania dla DAG-a. W Airflow DAG ma wymiar Data wykonania, tzn. w zależności od harmonogramu pracy DAG-a tworzone są instancje zadań dla każdej Daty wykonania. I dla każdej Daty wykonania zadania można wykonać ponownie — lub na przykład DAG może działać jednocześnie w kilku Datach wykonania. Jest to jasno pokazane tutaj:

Niestety (a może i na szczęście: zależy od sytuacji), jeśli zmienia się implementację zadania w DAG-u, to wykonanie w poprzednich Datach wykonania będzie już uwzględniać poprawki. To dobrze, jeśli trzeba przeliczyć dane w przeszłych okresach nowym algorytmem, ale źle, ponieważ traci się powtarzalność wyniku (oczywiście nikt nie przeszkadza w przywróceniu z Gita wymaganej wersji źródła i jednorazowym przeliczeniu tego, co potrzeba, tak, jak trzeba).
Generowanie zadań
Implementacja DAG-a to kod w Pythonie, dlatego mamy bardzo wygodny sposób na skrócenie objętości kodu podczas pracy, na przykład, z rozdzielonymi źródłami. Załóżmy, że masz trzy shardy MySQL jako źródło, musisz zajrzeć do każdego i pobrać jakieś dane. I to niezależnie i równolegle. Kod w Pythonie w DAG-u może wyglądać tak:
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 wygląda tak:

Można dodać lub usunąć shard, po prostu dostosowując ustawienia i aktualizując DAG. Wygodne!
Można także wykorzystać bardziej złożoną generację kodu, na przykład pracować z źródłami w postaci bazy danych lub opisywać strukturę tabeli, algorytm pracy z tabelą, a z uwzględnieniem specyfiki infrastruktury DWH generować proces ładowania N tabel do twojego magazynu. Możesz także generować N zadań w DAG-u z listy w API, które nie wspiera działania z parametrem w formie listy, ograniczyć równoległość zapytań do API puli i wydobyć z API potrzebne dane. Elastyczne!
Repozytorium
W Airflow istnieje własne repozytorium backendowe, baza danych (może to być MySQL lub Postgres, u nas Postgres), w której przechowywane są stany zadań, DAG-ów, ustawienia połączeń, zmienne globalne itd. Muszę zaznaczyć, że repozytorium w Airflow jest bardzo proste (około 20 tabel) i wygodne, jeśli chcesz zbudować jakiś swój proces na jego podstawie. Przypomina mi to 100500 tabel w repozytorium Informatica, które należało długo rozgryzać, zanim zrozumiano, jak stworzyć zapytanie.
Monitoring
Biorąc pod uwagę prostotę repozytorium, możesz samodzielnie zbudować wygodny dla siebie proces monitorowania zadań. Używamy notatnika w Zeppelin, gdzie obserwujemy stan zadań:

Może to być również interfejs webowy samego Airflow:

Kod Airflow jest otwarty, dlatego dodaliśmy alertowanie w Telegramie. Każdy działający instancja zadania, jeśli wystąpi błąd, wysyła powiadomienia do grupy na Telegramie, w której znajduje się cały zespół deweloperów i wsparcia.
Otrzymujemy przez Telegram szybkie reakcje (jeśli to konieczne), a przez Zeppelin - ogólny obraz zadań w Airflow.
Podsumowując
Airflow przede wszystkim jest open source i nie trzeba od niego oczekiwać cudów. Bądź gotowy na to, aby poświęcić czas i wysiłek na zbudowanie działającego rozwiązania. Cel jest osiągalny, uwierz mi, warto się o to postarać. Szybkość rozwoju, elastyczność, łatwość dodawania nowych procesów - spodoba ci się. Oczywiście, trzeba poświęcić dużo uwagi organizacji projektu i stabilności działania samego Airflow: cuda się nie zdarzają.
Obecnie nasz Airflow codziennie wykonuje około 6,5 tysiąca zadań. Są one dość różne. Istnieją zadania ładowania danych do głównego DWH z wielu różnych i bardzo specyficznych źródeł, są zadania obliczania witryn w ramach głównego DWH, są zadania publikacji danych w szybkim DWH, jest wiele - wiele różnych zadań - a Airflow to wszystko przetwarza dzień po dniu. Jeśli chodzi o liczby, to jest to 2,3 tysiąca złożonych zadań ELT wewnątrz DWH (Hadoop), około 2,5 setek baz danych źródeł, to zespół składa się z czterech programistów ETL, którzy dzielą się na procesowanie ETL danych w DWH oraz na procesowanie ELT danych wewnątrz DWH, a także na jednego administratora, który zajmuje się infrastrukturą usługi.
Plany na przyszłość
Liczba procesów nieuchronnie rośnie, a główną rzeczą, którą będziemy robić w zakresie infrastruktury Airflow, jest skalowanie. Chcemy zbudować klaster Airflow, przydzielić kilka nodów dla workerów Celery i stworzyć duplikującą się głowę z procesami planowania zadań oraz repozytorium.
Epilog
Oczywiście to nie wszystko, co chciałbym powiedzieć o Airflow, ale starałem się omówić najważniejsze punkty. Apetyt rośnie w miarę jedzenia, spróbuj - a spodoba ci się 🙂
Źródło: habr.com
