Airflow est un outil pour développer et maintenir facilement et rapidement des processus batch de traitement de données.

Airflow est un outil pour développer et maintenir facilement et rapidement des processus batch de traitement de données.

Bonjour, Habr ! Dans cet article, je souhaite vous parler d'un outil remarquable pour le dĂ©veloppement de processus batch de traitement de donnĂ©es, par exemple, dans l'infrastructure d'un DWH d'entreprise ou de votre DataLake. Nous allons parler d'Apache Airflow (ci-aprĂšs Airflow). Il est injustement sous-estimĂ© sur Habr, et dans la majeure partie de l'article, j'essaierai de vous convaincre qu'au moins Airflow mĂ©rite d'ĂȘtre considĂ©rĂ© lors du choix d'un planificateur pour vos processus ETL/ELT.

Auparavant, j'avais écrit une série d'articles sur le thÚme du DWH, lorsque je travaillais chez Tinkoff Bank. Maintenant, je fais partie de l'équipe de Mail.Ru Group et je m'occupe du développement d'une plateforme pour l'analyse de données dans le domaine du jeu. En fait, au fur et à mesure que des nouvelles et des solutions intéressantes apparaissent, mon équipe et moi allons parler ici de notre plateforme d'analyse de données.

Prologue

Alors, commençons. Qu'est-ce qu'Airflow ? C'est une bibliothĂšque (ou plutĂŽt un ensemble de bibliothĂšques) pour le dĂ©veloppement, la planification et le monitoring des workflows. La principale caractĂ©ristique d'Airflow : pour dĂ©crire (dĂ©velopper) les processus, on utilise du code en langage Python. D'oĂč dĂ©coule une multitude d'avantages pour l'organisation de votre projet et son dĂ©veloppement : en fait, votre projet ETL (par exemple) n'est rien d'autre qu'un projet Python, et vous pouvez l'organiser comme bon vous semble, en tenant compte des spĂ©cificitĂ©s de l'infrastructure, de la taille de l'Ă©quipe et d'autres exigences. CĂŽtĂ© outils, c'est simple. Utilisez, par exemple, PyCharm + Git. C'est gĂ©nial et trĂšs pratique !

Examinons maintenant les principales entités d'Airflow. En comprenant leur essence et leur fonction, vous optimiserez l'architecture de vos processus. Probablement, l'entité principale est le Directed Acyclic Graph (ci-aprÚs DAG).

DAG

Le DAG est une sorte d'union significative de vos tùches que vous souhaitez exécuter dans un ordre strictement défini selon un calendrier spécifique. Airflow offre une interface web conviviale pour travailler avec les DAG et d'autres entités :

Airflow est un outil pour développer et maintenir facilement et rapidement des processus batch de traitement de données.

Le DAG peut ressembler Ă  ceci :

Airflow est un outil pour développer et maintenir facilement et rapidement des processus batch de traitement de données.

Lorsqu'un développeur conçoit un DAG, il définit un ensemble d'opérateurs sur lesquels les tùches au sein du DAG seront basées. Nous arrivons ici à une autre entité importante : l'Airflow Operator.

Les opérateurs

L'OpĂ©rateur est une entitĂ© Ă  partir de laquelle des instances de tĂąches sont créées, dĂ©crivant ce qui se passera lors de l'exĂ©cution de l'instance de la tĂąche. Les versions d'Airflow depuis GitHub contiennent dĂ©jĂ  un ensemble d'opĂ©rateurs prĂȘts Ă  l'emploi. Exemples :

  • BashOperator — opĂ©rateur pour exĂ©cuter une commande bash.
  • PythonOperator — opĂ©rateur pour appeler du code Python.
  • EmailOperator — opĂ©rateur pour envoyer un email.
  • HTTPOperator — opĂ©rateur pour travailler avec des requĂȘtes http.
  • SqlOperator — opĂ©rateur pour exĂ©cuter du code SQL.
  • Sensor — opĂ©rateur d'attente d'un Ă©vĂ©nement (survenance d'un moment prĂ©cis, apparition d'un fichier requis, ligne dans une base de donnĂ©es, rĂ©ponse d'une API, etc.).

Il existe des opérateurs plus spécifiques : DockerOperator, HiveOperator, S3FileTransferOperator, PrestoToMysqlOperator, SlackOperator.

Vous pouvez Ă©galement dĂ©velopper des opĂ©rateurs en fonction de vos spĂ©cificitĂ©s et les utiliser dans le projet. Par exemple, nous avons créé MongoDBToHiveViaHdfsTransfer, un opĂ©rateur pour exporter des documents de MongoDB vers Hive, et plusieurs opĂ©rateurs pour travailler avec : CHLoadFromHiveOperator et CHTableLoaderOperator. En essence, dĂšs qu'une portion de code frĂ©quemment utilisĂ©e Ă©merge dans un projet, basĂ©e sur des opĂ©rateurs de base, on peut envisager de la rassembler en un nouvel opĂ©rateur. Cela simplifiera le dĂ©veloppement futur et enrichira votre bibliothĂšque d'opĂ©rateurs dans le projet. ClickHouseEnsuite, tous ces instances de tĂąches doivent ĂȘtre exĂ©cutĂ©es, et parlons maintenant du planificateur.

Le planificateur de tùches dans Airflow est basé sur

Planificateur

Celery. Celery est une bibliothĂšque Python qui permet d'organiser une file d'attente ainsi qu'une exĂ©cution asynchrone et distribuĂ©e des tĂąches. Du cĂŽtĂ© d'Airflow, toutes les tĂąches sont rĂ©parties en pools. Les pools sont créés manuellement. Leur objectif est gĂ©nĂ©ralement de limiter la charge sur le travail avec la source ou de typifier les tĂąches Ă  l'intĂ©rieur du DWH. Les pools peuvent ĂȘtre gĂ©rĂ©s via l'interface web :Chaque pool a une limite de nombre de slots. Lors de la crĂ©ation d'un DAG, il se voit assigner un pool :

Airflow est un outil pour développer et maintenir facilement et rapidement des processus batch de traitement de données.

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__

Le pool dĂ©fini au niveau du DAG peut ĂȘtre redĂ©fini au niveau de la tĂąche.

La planification de toutes les tĂąches dans Airflow est assurĂ©e par un processus distinct — le Scheduler. En rĂ©alitĂ©, le Scheduler s'occupe de toute la mĂ©canique de la mise en exĂ©cution des tĂąches. Avant d'ĂȘtre exĂ©cutĂ©e, une tĂąche passe par plusieurs Ă©tapes :
La planification de toutes les tĂąches dans Airflow est gĂ©rĂ©e par un processus distinct : le Scheduler. En effet, le Scheduler s'occupe de toute la mĂ©canique de mise en Ɠuvre des tĂąches. Avant qu'une tĂąche ne soit exĂ©cutĂ©e, elle passe par plusieurs Ă©tapes :

  1. Dans le DAG, les tĂąches prĂ©cĂ©dentes ont Ă©tĂ© exĂ©cutĂ©es, une nouvelle peut ĂȘtre mise en file d'attente.
  2. La file d'attente est triĂ©e en fonction de la prioritĂ© des tĂąches (ces prioritĂ©s peuvent Ă©galement ĂȘtre gĂ©rĂ©es), et si un slot est libre dans le pool, la tĂąche peut ĂȘtre prise en charge.
  3. S'il y a un worker celery libre, la tùche lui est envoyée ; le travail que vous avez programmé dans la tùche commence, en utilisant un opérateur approprié.

C'est assez simple.

Le Scheduler fonctionne sur tous les DAG et toutes les tùches à l'intérieur des DAG.

Pour que le Scheduler commence Ă  travailler avec un DAG, il faut lui attribuer un calendrier :

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

Il existe un ensemble de presets prĂȘts Ă  l'emploi : @once, @hourly, @daily, @weekly, @monthly, @yearly.

On peut également utiliser des expressions cron :

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

Date d'exécution

Pour comprendre comment Airflow fonctionne, il est important de saisir ce qu'est la Date d'exĂ©cution pour un DAG. Dans Airflow, un DAG a une dimension Date d'exĂ©cution, c'est-Ă -dire que des instances de tĂąches sont créées pour chaque Date d'exĂ©cution en fonction du calendrier de travail du DAG. Et pour chaque Date d'exĂ©cution, les tĂąches peuvent ĂȘtre exĂ©cutĂ©es Ă  nouveau — ou, par exemple, le DAG peut fonctionner simultanĂ©ment sur plusieurs Dates d'exĂ©cution. Cela est clairement illustrĂ© ici :

Airflow est un outil pour développer et maintenir facilement et rapidement des processus batch de traitement de données.

Malheureusement (ou peut-ĂȘtre heureusement : cela dĂ©pend de la situation), si la mise en Ɠuvre de la tĂąche dans le DAG est modifiĂ©e, son exĂ©cution dans les Dates d'exĂ©cution prĂ©cĂ©dentes se fera dĂ©jĂ  en tenant compte des ajustements. C'est bien si vous avez besoin de recalculer des donnĂ©es pour des pĂ©riodes passĂ©es avec un nouvel algorithme, mais c'est mal car cela compromet la reproductibilitĂ© du rĂ©sultat (bien sĂ»r, rien n'empĂȘche de revenir Ă  Git pour trouver la version souhaitĂ©e du code source et de recalculer ce qui est nĂ©cessaire, comme il le faut).

Génération de tùches

La mise en Ɠuvre d'un DAG est du code Python, c'est pourquoi nous avons un moyen trĂšs pratique de rĂ©duire le volume de code lorsque nous travaillons, par exemple, avec des sources shardĂ©es. Supposons que vous ayez trois shards MySQL en tant que source ; vous devez vous connecter Ă  chacun d'eux et rĂ©cupĂ©rer certaines donnĂ©es. De maniĂšre indĂ©pendante et en parallĂšle. Le code Python dans le DAG peut ressembler Ă  ceci :

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)

Le DAG est ainsi constitué :

Airflow est un outil pour développer et maintenir facilement et rapidement des processus batch de traitement de données.

Vous pouvez ajouter ou supprimer un shard en ajustant simplement les paramĂštres et en mettant Ă  jour le DAG. Pratique!

Il est possible d'utiliser une gĂ©nĂ©ration de code plus complexe, par exemple en travaillant avec des sources sous forme de bases de donnĂ©es ou en dĂ©crivant une structure tabulaire, un algorithme de traitement de tableau, et en gĂ©nĂ©rant un processus de chargement de N tables dans votre entrepĂŽt en tenant compte des particularitĂ©s de l'infrastructure DWH. Ou, par exemple, en travaillant avec une API qui ne prend pas en charge le traitement de paramĂštres sous forme de liste, vous pouvez gĂ©nĂ©rer N tĂąches dans le DAG Ă  partir de cette liste, limiter la parallĂ©lisation des requĂȘtes API avec un pool et extraire les donnĂ©es nĂ©cessaires de l'API. Flexible!

DépÎt

Airflow dispose de son propre rĂ©fĂ©rentiel backend, une base de donnĂ©es (qui peut ĂȘtre MySQL ou Postgres, chez nous c'est Postgres), oĂč sont stockĂ©s l'Ă©tat des tĂąches, les DAG, les configurations de connexions, les variables globales, etc. Il convient de mentionner que le rĂ©fĂ©rentiel d'Airflow est trĂšs simple (environ 20 tables) et pratique si vous souhaitez construire un processus propre. Cela me rappelle les 100 500 tables dans le rĂ©fĂ©rentiel Informatica, qu'il fallait longtemps Ă©tudier avant de comprendre comment formuler une requĂȘte.

Surveillance

Étant donnĂ© la simplicitĂ© du rĂ©fĂ©rentiel, vous pouvez construire vous-mĂȘme un processus de surveillance des tĂąches qui vous convient. Nous utilisons un carnet de notes dans Zeppelin, oĂč nous vĂ©rifions l'Ă©tat des tĂąches :

Airflow est un outil pour développer et maintenir facilement et rapidement des processus batch de traitement de données.

Cela peut Ă©galement ĂȘtre l'interface web d'Airflow :

Airflow est un outil pour développer et maintenir facilement et rapidement des processus batch de traitement de données.

Le code d'Airflow est ouvert, donc nous avons ajoutĂ© des alertes dans Telegram. Chaque instance de tĂąche en cours, en cas d'erreur, envoie un message dans le groupe Telegram oĂč se trouve toute l'Ă©quipe de dĂ©veloppement et de support.

Nous recevons une réaction rapide via Telegram (si nécessaire), et une vue d'ensemble des tùches dans Airflow via Zeppelin.

Au total

Airflow est avant tout open source, et il ne faut pas s'attendre Ă  des miracles. Soyez prĂȘts Ă  consacrer du temps et des efforts pour mettre en place une solution fonctionnelle. L'objectif est atteignable, croyez-moi, cela en vaut la peine. La rapiditĂ© de dĂ©veloppement, la flexibilitĂ©, la simplicitĂ© d'ajout de nouveaux processus — cela vous plaira. Bien sĂ»r, il est nĂ©cessaire de prĂȘter beaucoup d'attention Ă  l'organisation du projet et Ă  la stabilitĂ© de fonctionnement d'Airflow : il n'y a pas de miracles.

Actuellement, nous avons Airflow qui exĂ©cute quotidiennement environ 6 500 tĂąches.. Leur nature est assez diffĂ©rente. Il y a des tĂąches de chargement de donnĂ©es vers l'entrepĂŽt de donnĂ©es principal Ă  partir de nombreuses sources variĂ©es et trĂšs spĂ©cifiques, des tĂąches de calcul de vitrines Ă  l'intĂ©rieur de l'entrepĂŽt de donnĂ©es principal, des tĂąches de publication de donnĂ©es dans un entrepĂŽt de donnĂ©es rapide, et beaucoup, beaucoup d'autres tĂąches — et Airflow les gĂšre jour aprĂšs jour. Si l'on parle en chiffres, c'est 2,3 mille des tĂąches ELT de complexitĂ© variĂ©e Ă  l'intĂ©rieur de l'entrepĂŽt de donnĂ©es (Hadoop), environ 250 bases de donnĂ©es sources, c'est une Ă©quipe de quatre dĂ©veloppeurs ETL, qui se rĂ©partissent entre le traitement ETL des donnĂ©es dans l'entrepĂŽt de donnĂ©es et le traitement ELT des donnĂ©es Ă  l'intĂ©rieur de l'entrepĂŽt de donnĂ©es, et bien sĂ»r encore un administrateur, qui s'occupe de l'infrastructure du service.

Plans pour l'avenir

Le nombre de processus ne peut que croĂźtre, et ce sur quoi nous allons nous concentrer en matiĂšre d'infrastructure Airflow, c'est l'Ă©chelle. Nous voulons construire un cluster Airflow, allouer quelques nƓuds pour les workers Celery et crĂ©er une tĂȘte redondante avec des processus de planification de tĂąches et un rĂ©fĂ©rentiel.

Épilogue

Ce n'est bien sĂ»r pas tout ce que j'aimerais partager sur Airflow, mais j'ai essayĂ© de couvrir les points principaux. L'appĂ©tit vient en mangeant, essayez — et vous aimerez 🙂

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