Apache Airflow : simplifions l'ETL

Bonjour, je suis Dmitri Logvinenko, Data Engineer au sein du département d'analyse du groupe «Vezet».

Je vais vous parler d'un outil remarquable pour le dĂ©veloppement de processus ETL : Apache Airflow. Mais Airflow est si polyvalent et riche en fonctionnalitĂ©s qu'il vaut la peine de s'y intĂ©resser mĂȘme si vous ne travaillez pas avec des flux de donnĂ©es, mais avez besoin de lancer pĂ©riodiquement certains processus et de suivre leur exĂ©cution.

Et oui, je vais non seulement vous expliquer, mais aussi vous montrer : le programme contient beaucoup de code, de captures d'écran et de recommandations.

Apache Airflow : simplifions l'ETL
Que voit-on habituellement lorsque l'on recherche le mot Airflow / Wikimedia Commons

Table des matiĂšres

Introduction

Apache Airflow — c'est un peu comme Django :

  • Ă©crit en Python,
  • il offre une excellente interface admin,
  • il est extensible Ă  l'infini,

— juste mieux, et conçu pour des objectifs complĂštement diffĂ©rents, Ă  savoir (comme Ă©crit jusqu'ici) :

  • l'exĂ©cution et la surveillance de tĂąches sur un nombre illimitĂ© de machines (autant que Celery/Kubernetes et votre conscience le permettent)
  • avec une gĂ©nĂ©ration dynamique de workflow Ă  partir d'un code Python trĂšs simple Ă  Ă©crire et Ă  comprendre
  • et la possibilitĂ© de relier n'importe quelles bases de donnĂ©es et API, Ă  l'aide Ă  la fois de composants prĂȘts Ă  l'emploi et de plugins faits maison (ce qui est extrĂȘmement simple Ă  rĂ©aliser).

Nous utilisons Apache Airflow de la maniĂšre suivante :

  • nous collectons des donnĂ©es provenant de diffĂ©rentes sources (de nombreux instances SQL Server et PostgreSQL, diverses API avec les mĂ©triques des applications, mĂȘme 1C) dans DWH et ODS (nous utilisons Vertica et Clickhouse).
  • comme un systĂšme avancĂ© cron, qui lance des processus de consolidation des donnĂ©es dans ODS, tout en surveillant leur maintenance.

Jusqu'Ă  rĂ©cemment, nos besoins Ă©taient couverts par un petit serveur de 32 cƓurs et 50 Go de RAM. Dans Airflow, il fonctionne avec :

  • plus de 200 DAGs (en fait des workflows dans lesquels nous avons inclus des tĂąches),
  • avec en moyenne 70 tĂąches,
  • ceci s'exĂ©cutant (Ă©galement en moyenne) une fois par heure.

Et je vais parler de notre expansion ci-dessous, mais pour l'instant, dĂ©finissons la tĂąche ĂŒber que nous allons rĂ©soudre :

Il existe trois SQL Servers d'origine, chacun avec 50 bases de donnĂ©es — des instances d'un mĂȘme projet, donc leur structure est presque identique (partout, muahaha), ce qui signifie qu'il y a une table Orders dans chacune (heureusement, une table avec ce nom peut ĂȘtre intĂ©grĂ©e dans n'importe quel secteur d'activitĂ©). Nous rĂ©cupĂ©rons les donnĂ©es en ajoutant des champs de mĂ©tadonnĂ©es (serveur source, base source, identifiant de la tĂąche ETL) et naĂŻvement, nous allons les jeter dans, disons, Vertica.

Allons-y !

Partie principale, pratique (et un peu théorique)

Pourquoi est-ce intéressant pour nous (et pour vous)

Quand les arbres étaient grands et que j'étais simple SQL-moulins dans un grand détaillant russe, nous pilotions les processus ETL aka les flux de données à l'aide de deux outils que nous avions à disposition :

  • Informatica Power Center — un systĂšme extrĂȘmement complexe et trĂšs performant, avec son propre matĂ©riel et son propre systĂšme de versioning. Je n'ai utilisĂ© mĂȘme pas 1% de ses capacitĂ©s. Pourquoi ? Eh bien, d'une part, cette interface vient d'un autre temps et nous pesait psychologiquement. D'autre part, cette machine est conçue pour des processus extrĂȘmement sophistiquĂ©s, une rĂ©utilisation fĂ©roce des composants et d'autres fonctionnalitĂ©s trĂšs importantes pour les entreprises. En ce qui concerne son coĂ»t, comparable Ă  celui d'une aile d'Airbus A380 par an, nous allons rester silencieux.

    Attention, une capture d'écran peut blesser un peu les personnes de moins de 30 ans

    Apache Airflow : simplifions l'ETL

  • SQL Server Integration Server — nous avons utilisĂ© ce camarade dans nos flux internes de projet. En rĂ©alitĂ© : nous utilisons dĂ©jĂ  SQL Server, et il serait un peu irrationnel de ne pas utiliser ses outils ETL. Tout dans celui-ci est agrĂ©able : l'interface est belle, et les rapports d'exĂ©cution... Mais ce n'est pas pour cela que nous aimons les produits logiciels, oh non. Nous pouvons versionner ses dtsx (qui est un XML avec des nƓuds mĂ©langĂ©s lors de la sauvegarde), mais Ă  quoi bon ? Et crĂ©er un paquet de tĂąches qui transfĂ©rera une centaine de tables d'un serveur Ă  un autre ? Pff, qu'une centaine, aprĂšs une vingtaine de clics, votre index va flancher sur le bouton de la souris. Mais il a, sans aucun doute, un look plus Ă  la mode :

    Apache Airflow : simplifions l'ETL

Nous cherchions sans cesse des solutions. Le fait est que presque nous en sommes arrivés à un générateur de paquets SSIS fait maison...

... puis un nouvel emploi m'a trouvé. Et dans ce nouvel emploi, Apache Airflow m'a rattrapé.

Quand j'ai appris que la description des processus ETL était simplement du code Python, je ne pouvais presque pas contenir ma joie. Ainsi, les flux de données ont été soumis à la version et à la différenciation, et rassembler des tables avec une structure unique à partir de centaines de bases de données en une seule cible est devenu une question de quelques lignes de code Python sur un écran de 13 pouces.

Assembler un cluster

Évitons de faire un vĂ©ritable jardin d'enfants et parlons ici de choses totalement Ă©videntes, comme l'installation d'Airflow, la base de donnĂ©es que vous avez choisie, Celery et d'autres aspects dĂ©crits dans la documentation.

Afin que nous puissions commencer immédiatement les expérimentations, j'ai esquissé docker-compose.yml where:

  • Levant notre propre Airflow: Scheduler, Webserver. Flower sera Ă©galement utilisĂ© pour surveiller les tĂąches Celery (car il a dĂ©jĂ  Ă©tĂ© intĂ©grĂ© dans apache/airflow:1.10.10-python3.7, et nous n'y voyons pas d'inconvĂ©nient);
  • PostgreSQL, oĂč Airflow Ă©crira ses informations de service (donnĂ©es du planificateur, statistiques d'exĂ©cution, etc.), et Celery - marquera les tĂąches terminĂ©es;
  • Redis, qui agira comme un courtier de tĂąches pour Celery;
  • Celery worker, qui s'occupera de l'exĂ©cution des tĂąches.
  • Dans le dossier . /dags nous allons placer nos fichiers dĂ©crivant les dags. Ils seront rĂ©cupĂ©rĂ©s Ă  la volĂ©e, donc il n'est pas nĂ©cessaire de redĂ©marrer toute la pile aprĂšs chaque petit changement.

Dans certains endroits, le code dans les exemples n'est pas complĂštement prĂ©sentĂ© (pour ne pas alourdir le texte), et ailleurs, il est modifiĂ© en cours de route. Des exemples de code complets et fonctionnels peuvent ĂȘtre consultĂ©s dans le dĂ©pĂŽt. https://github.com/dm-logv/airflow-tutorial.

docker-compose.yml

version: '3.4'

x-airflow-config: &airflow-config
  AIRFLOW__CORE__DAGS_FOLDER: /dags
  AIRFLOW__CORE__EXECUTOR: CeleryExecutor
  AIRFLOW__CORE__FERNET_KEY: MJNz36Q8222VOQhBOmBROFrmeSxNOgTCMaVp2_HOtE0=
  AIRFLOW__CORE__HOSTNAME_CALLABLE: airflow.utils.net:get_host_ip_address
  AIRFLOW__CORE__SQL_ALCHEMY_CONN: postgres+psycopg2://airflow:airflow@airflow-db:5432/airflow

  AIRFLOW__CORE__PARALLELISM: 128
  AIRFLOW__CORE__DAG_CONCURRENCY: 16
  AIRFLOW__CORE__MAX_ACTIVE_RUNS_PER_DAG: 4
  AIRFLOW__CORE__LOAD_EXAMPLES: 'False'
  AIRFLOW__CORE__LOAD_DEFAULT_CONNECTIONS: 'False'

  AIRFLOW__EMAIL__DEFAULT_EMAIL_ON_RETRY: 'False'
  AIRFLOW__EMAIL__DEFAULT_EMAIL_ON_FAILURE: 'False'

  AIRFLOW__CELERY__BROKER_URL: redis://broker:6379/0
  AIRFLOW__CELERY__RESULT_BACKEND: db+postgresql://airflow:airflow@airflow-db/airflow

x-airflow-base: &airflow-base
  image: apache/airflow:1.10.10-python3.7
  entrypoint: /bin/bash
  restart: always
  volumes:
    - ./dags:/dags
    - ./requirements.txt:/requirements.txt

services:
  # Redis comme broker Celery
  broker:
    image: redis:6.0.5-alpine

  # Base de données pour les métadonnées Airflow
  airflow-db:
    image: postgres:10.13-alpine

    environment:
      - POSTGRES_USER=airflow
      - POSTGRES_PASSWORD=airflow
      - POSTGRES_DB=airflow

    volumes:
      - ./db:/var/lib/postgresql/data

  # Conteneur principal avec Airflow Webserver, Scheduler, Celery Flower
  airflow:
    <<: *airflow-base

    environment:
      <<: *airflow-config

      AIRFLOW__SCHEDULER__DAG_DIR_LIST_INTERVAL: 30
      AIRFLOW__SCHEDULER__CATCHUP_BY_DEFAULT: 'False'
      AIRFLOW__SCHEDULER__MAX_THREADS: 8

      AIRFLOW__WEBSERVER__LOG_FETCH_TIMEOUT_SEC: 10

    depends_on:
      - airflow-db
      - broker

    command: >
      -c " sleep 10 &&
           pip install --user -r /requirements.txt &&
           /entrypoint initdb &&
          (/entrypoint webserver &) &&
          (/entrypoint flower &) &&
           /entrypoint scheduler"

    ports:
      # Celery Flower
      - 5555:5555
      # Airflow Webserver
      - 8080:8080

  # Celery worker, sera redimensionné en utilisant `--scale=n`
  worker:
    <<: *airflow-base

    environment:
      <<: *airflow-config

    command: >
      -c " sleep 10 &&
           pip install --user -r /requirements.txt &&
           /entrypoint worker"

    depends_on:
      - airflow
      - airflow-db
      - broker

Remarques :

  • Dans la construction du compose, je me suis largement basĂ© sur l'image bien connue puckel/docker-airflow – Ă  ne pas manquer. Peut-ĂȘtre que vous n'aurez besoin de rien d'autre dans votre vie.
  • Tous les paramĂštres d'Airflow sont accessibles non seulement par airflow.cfg, mais aussi par des variables d'environnement (merci aux dĂ©veloppeurs), ce dont j'ai profitĂ© avec malice.
  • Bien sĂ»r, ce n'est pas prĂȘt pour la production : je n'ai pas dĂ©ployĂ© de heartbeats sur les conteneurs et je n'ai pas pris en compte la sĂ©curitĂ©. Mais j'ai créé un minimum adaptĂ© pour nos expĂ©rimentations.
  • Notez que :
    • Le dossier contenant les DAG doit ĂȘtre accessible Ă  la fois par le planificateur et par les travailleurs.
    • Il en va de mĂȘme pour toutes les bibliothĂšques tierces — elles doivent toutes ĂȘtre installĂ©es sur les machines avec le planificateur et les travailleurs.

Eh bien, maintenant c'est simple :

$ docker-compose up --scale worker=3

AprĂšs que tout s'est mis en route, vous pouvez consulter les interfaces web :

Concepts de base

Si vous n'avez rien compris à tous ces « DAGs », voici un petit glossaire :

  • Scheduler — le gros bonhomme d'Airflow, qui veille Ă  ce que ce soient des robots qui bossent, pas des humains : il surveille le planning, met Ă  jour les DAGs, lance les tĂąches.

    En fait, dans les anciennes versions, il avait des problĂšmes de mĂ©moire (non, pas d'amnĂ©sie, mais des fuites) et il restait mĂȘme un paramĂštre legacy dans les configs run_duration — l'intervalle de son redĂ©marrage. Mais maintenant, tout va bien.

  • DAG (Ă©galement appelĂ© « dag ») — « graphe acyclique orientĂ© », mais cette dĂ©finition ne dira pas grand-chose Ă  la plupart des gens, en rĂ©alitĂ© c'est un conteneur pour des tĂąches qui interagissent entre elles (voir ci-dessous) ou l'analogue de Package dans SSIS et Workflow dans Informatica.

    En plus des DAGs, il peut y avoir des sous-DAGs, mais nous n'y arriverons probablement pas.

  • DAG Run — un DAG initialisĂ© auquel est attribuĂ©e sa propre execution_date. Les DAG runs d'un mĂȘme DAG peuvent parfaitement fonctionner en parallĂšle (si bien sĂ»r vous avez rendu vos tĂąches idempotentes).
  • Operator — ce sont des morceaux de code responsables de l'exĂ©cution d'une action spĂ©cifique. Il existe trois types d'opĂ©rateurs :
    • action, comme notre cher PythonOperator, qui peut exĂ©cuter n'importe quel code Python (valide) ;
    • transfert, qui dĂ©placent des donnĂ©es d'un endroit Ă  un autre, par exemple, MsSqlToHiveTransfer;
    • sensor permettra de rĂ©agir ou de ralentir l'exĂ©cution du DAG jusqu'Ă  ce qu'un Ă©vĂ©nement se produise. HttpSensor peut interroger un endpoint spĂ©cifiĂ©, et lorsque la bonne rĂ©ponse est reçue, lancer le transfert GoogleCloudStorageToS3Operator. Un esprit curieux pourrait demander : « pourquoi ? AprĂšs tout, on pourrait faire des rĂ©pĂ©titions directement dans l'opĂ©rateur ! » Et bien, pour ne pas saturer le pool de tĂąches avec des opĂ©rateurs suspendus. Le capteur se dĂ©clenche, vĂ©rifie et s'arrĂȘte jusqu'Ă  la prochaine tentative.
  • Task — les opĂ©rateurs dĂ©clarĂ©s, quel que soit leur type, attachĂ©s au DAG sont Ă©levĂ©s au rang de tĂąche.
  • Task instance — quand le planificateur gĂ©nĂ©ral a dĂ©cidĂ© qu'il Ă©tait temps d'envoyer les tĂąches au combat pour les exĂ©cuteurs-travailleurs (directement sur place, si nous utilisons LocalExecutor ou sur un nƓud distant dans le cas de CeleryExecutor), il leur attribue un contexte (c'est-Ă -dire un ensemble de variables — paramĂštres d'exĂ©cution), dĂ©ploie les modĂšles de commandes ou de requĂȘtes et les regroupe dans un pool.

Générer des tùches

Commençons par définir le schéma général de notre DAG, puis nous plongerons progressivement dans les détails, car nous appliquons certaines solutions non triviales.

Ainsi, dans sa forme la plus simple, un tel DAG ressemblera Ă  ceci :

from datetime import timedelta, datetime

from airflow import DAG
from airflow.operators.python_operator import PythonOperator

from commons.datasources import sql_server_ds

dag = DAG('orders',
          schedule_interval=timedelta(hours=6),
          start_date=datetime(2020, 7, 8, 0))

def workflow(**context):
    print(context)

for conn_id, schema in sql_server_ds:
    PythonOperator(
        task_id=schema,
        python_callable=workflow,
        provide_context=True,
        dag=dag)

Commençons par comprendre :

  • Tout d'abord, importons les bibliothĂšques nĂ©cessaires et quelques autres Ă©lĂ©ments;
  • sql_server_ds — ce sont des List[namedtuple[str, str]] avec les noms de connexions provenant des connexions Airflow et les bases de donnĂ©es d'oĂč nous allons rĂ©cupĂ©rer notre tableau ;
  • dag — la dĂ©claration de notre dag, qui doit absolument se trouver dans globals(), sinon Airflow ne le trouvera pas. Il faut Ă©galement dire au dag :
    • quel est son nom orders — ce nom apparaĂźtra ensuite dans l'interface web,
    • qu'il commencera Ă  fonctionner Ă  partir de minuit le 8 juillet,
    • et qu'il doit se dĂ©clencher environ toutes les 6 heures (pour les durs ici, Ă  la place de timedelta() on peut utiliser une cron-chaĂźne 0 0 0/6 ? * * *, pour les moins durs — une expression comme @daily);
  • workflow() fera le travail principal, mais pas maintenant. Pour l'instant, nous allons simplement afficher notre contexte dans le journal.
  • Et maintenant, la simple magie de crĂ©ation des tĂąches :
    • nous parcourons nos sources ;
    • nous initialisons PythonOperator, qui exĂ©cutera notre tĂąche vide workflow(). N'oubliez pas d'indiquer un nom unique (dans le cadre du dag) pour la tĂąche et d'attacher le dag lui-mĂȘme. Le flag provide_context injectera Ă  son tour des arguments supplĂ©mentaires dans la fonction, que nous rassemblerons avec soin grĂące Ă  **context.

Pour l'instant, c'est tout. Que nous avons obtenu :

  • un nouveau dag dans l'interface web,
  • une centaine de tĂąches qui seront exĂ©cutĂ©es en parallĂšle (si les paramĂštres d'Airflow, Celery et la puissance des serveurs le permettent).

Eh bien, presque obtenu.

Apache Airflow : simplifions l'ETL
Qui va établir les dépendances ?

Pour simplifier tout cela, j'ai intĂ©grĂ© dans docker-compose.yml le traitement requirements.txt sur tous les nƓuds.

Maintenant, ça y est :

Apache Airflow : simplifions l'ETL

Les carrĂ©s gris — instances de tĂąches, traitĂ©es par le planificateur.

Nous attendons un peu, les tĂąches sont prises par les workers :

Apache Airflow : simplifions l'ETL

Les vertes, c'est Ă©vident, — celles qui ont Ă©tĂ© exĂ©cutĂ©es avec succĂšs. Les rouges — pas trĂšs rĂ©ussies.

D'ailleurs, sur notre prod, il n'y a aucun dossier . /dags, synchronisĂ© entre les machines — tous les dags se trouvent dans git notre Gitlab, et Gitlab CI dĂ©ploie les mises Ă  jour sur les machines lors d'un merge dans master.

Un mot sur Flower

Pendant que les workers traitent nos tñches vides, rappelons-nous d'un autre outil, qui peut nous montrer certaines choses — Flower.

La premiĂšre page avec un rĂ©sumĂ© des informations sur les nƓuds-workers :

Apache Airflow : simplifions l'ETL

La page la plus chargée avec les tùches envoyées au travail :

Apache Airflow : simplifions l'ETL

La page la plus ennuyeuse sur l'état de notre courtier :

Apache Airflow : simplifions l'ETL

La page la plus dynamique — avec des graphiques de l'Ă©tat des tĂąches et leurs temps d'exĂ©cution :

Apache Airflow : simplifions l'ETL

Charger ce qui n'est pas encore chargé

Ainsi, toutes les tùches ont été effectuées, nous pouvons évacuer les blessés.

Apache Airflow : simplifions l'ETL

Et il y avait pas mal de blessĂ©s — pour diverses raisons. En cas d'utilisation correcte d'Airflow, ces carrĂ©s montrent que les donnĂ©es n'ont clairement pas Ă©tĂ© livrĂ©es.

Il faut consulter le journal et redémarrer les instances de tùche échouées.

En cliquant sur n'importe quel carré, nous verrons les actions qui nous sont disponibles :

Apache Airflow : simplifions l'ETL

On peut choisir de faire Clear sur celui qui a Ă©chouĂ©. Autrement dit, nous oublions que quelque chose s'est enrayĂ©, et la mĂȘme instance de tĂąche partira vers le planificateur.

Apache Airflow : simplifions l'ETL

Évidemment, procĂ©der ainsi avec toutes les cases rouges n'est pas trĂšs humain — ce n'est pas ce que nous attendons d'Airflow. Naturellement, nous avons une arme de destruction massive : Browse/Task Instances

Apache Airflow : simplifions l'ETL

Sélectionnons tout d'un coup et appliquons la bonne option :

Apache Airflow : simplifions l'ETL

AprÚs nettoyage, nos tùches ressemblent à ceci (elles attendent déjà avec impatience que le planificateur les programme) :

Apache Airflow : simplifions l'ETL

Connexions, hooks et autres variables

Il est temps de jeter un Ɠil au prochain DAG, update_reports.py:

from collections import namedtuple
from datetime import datetime, timedelta
from textwrap import dedent

from airflow import DAG
from airflow.contrib.operators.vertica_operator import VerticaOperator
from airflow.operators.email_operator import EmailOperator
from airflow.utils.trigger_rule import TriggerRule

from commons.operators import TelegramBotSendMessage

dag = DAG('update_reports',
          start_date=datetime(2020, 6, 7, 6),
          schedule_interval=timedelta(days=1),
          default_args={'retries': 3, 'retry_delay': timedelta(seconds=10)})

Report = namedtuple('Report', 'source target')
reports = [Report(f'{table}_view', table) for table in [
    'reports.city_orders',
    'reports.client_calls',
    'reports.client_rates',
    'reports.daily_orders',
    'reports.order_duration']]

email = EmailOperator(
    task_id='email_success', dag=dag,
    to='{{ var.value.all_the_kings_men }}',
    subject='DWH Reports updated',
    html_content=dedent("""Mesdames et messieurs, les rapports ont été mis à jour"""),
    trigger_rule=TriggerRule.ALL_SUCCESS)

tg = TelegramBotSendMessage(
    task_id='telegram_fail', dag=dag,
    tg_bot_conn_id='tg_main',
    chat_id='{{ var.value.failures_chat }}',
    message=dedent("""
         Natasha, réveille-toi, nous avons fait tomber {{ dag.dag_id }}
        """),
    trigger_rule=TriggerRule.ONE_FAILED)

for source, target in reports:
    queries = [f"TRUNCATE TABLE {target}",
               f"INSERT INTO {target} SELECT * FROM {source}"]

    report_update = VerticaOperator(
        task_id=target.replace('reports.', ''),
        sql=queries, vertica_conn_id='dwh',
        task_concurrency=1, dag=dag)

    report_update >> [email, tg]

Tout le monde a-t-il déjà fait un mise à jour de rapports ? C'est encore elle : il y a une liste de sources à partir desquelles récupérer les données ; il y a une liste d'endroits pour les placer ; n'oublions pas de signaler lorsque tout s'est passé ou s'est cassé (bon, ça ne nous concerne pas, non).

Reprenons le fichier et examinons les nouveaux éléments curieux :

  • from commons.operators import TelegramBotSendMessage — rien ne nous empĂȘche de crĂ©er nos propres opĂ©rateurs, ce que nous avons fait en rĂ©alisant un petit wrapper pour envoyer des messages dans Unblocked. (Nous parlerons plus en dĂ©tail de cet opĂ©rateur ci-dessous);
  • default_args={} — le DAG peut distribuer les mĂȘmes arguments Ă  tous ses opĂ©rateurs;
  • to='{{ var.value.all_the_kings_men }}' — champ to nous ne le coderons pas en dur, mais nous le gĂ©nĂ©rerons dynamiquement Ă  l'aide de Jinja et d'une variable contenant une liste d'adresses email, que j'ai soigneusement placĂ©e dans Admin/Variables;
  • trigger_rule=TriggerRule.ALL_SUCCESS — condition de dĂ©clenchement de l'opĂ©rateur. Dans notre cas, l'email sera envoyĂ© aux chefs seulement si toutes les dĂ©pendances se sont exĂ©cutĂ©es avec succĂšs;
  • tg_bot_conn_id='tg_main' — arguments conn_id reçoivent les identifiants de connexions que nous crĂ©ons dans Admin/Connections;
  • trigger_rule=TriggerRule.ONE_FAILED — les messages dans Telegram seront envoyĂ©s uniquement en cas d'Ă©chec de tĂąches;
  • task_concurrency=1 — nous interdisons le lancement simultanĂ© de plusieurs instances de tĂąches du mĂȘme task. Sinon, nous obtiendrons le lancement simultanĂ© de plusieurs VerticaOperator (regardant la mĂȘme table);
  • report_update >> [email, tg] — tout VerticaOperator se rejoindra pour l'envoi de l'email et du message, comme ceci :
    Apache Airflow : simplifions l'ETL

    Mais comme les opérateurs de notification ont des conditions de déclenchement différentes, seul un d'entre eux fonctionnera. Dans l'arborescence, tout semble un peu moins clair :
    Apache Airflow : simplifions l'ETL

Je vais dire quelques mots sur les macros et leurs amis — des variables.

Les macros sont des placeholders Jinja qui peuvent insérer diverses informations utiles dans les arguments des opérateurs. Par exemple, comme ceci :

SELECT
    id,
    payment_dtm,
    payment_type,
    client_id
FROM orders.payments
WHERE
    payment_dtm::DATE = '{{ ds }}'::DATE

{{ ds }} se dĂ©veloppera en le contenu de la variable de contexte execution_date au format YYYY-MM-DD: 2020-07-14. Ce qui est agrĂ©able, c'est que les variables de contexte sont attachĂ©es Ă  une instance de tĂąche donnĂ©e (la case dans l'arborescence), et lors du redĂ©marrage, les placeholders seront dĂ©veloppĂ©s dans les mĂȘmes valeurs.

Les valeurs assignĂ©es peuvent ĂȘtre consultĂ©es Ă  l'aide du bouton Rendered sur chaque instance de tĂąche. Voici comment cela se prĂ©sente pour la tĂąche d'envoi d'email :

Apache Airflow : simplifions l'ETL

Et voici pour la tĂąche d'envoi de message :

Apache Airflow : simplifions l'ETL

Une liste complÚte des macros intégrées pour la derniÚre version disponible est disponible ici : Macros Reference

De plus, grùce aux plugins, nous pouvons déclarer nos propres macros, mais c'est une toute autre histoire.

En plus des éléments prédéfinis, nous pouvons injecter des valeurs de nos variables (j'ai déjà utilisé cela plus haut dans le code). Créons dans Admin/Variables quelques éléments :

Apache Airflow : simplifions l'ETL

VoilĂ , il est prĂȘt Ă  l'emploi :

TelegramBotSendMessage(chat_id='{{ var.value.failures_chat }}')

Dans la valeur, il peut y avoir un scalaire, ou mĂȘme un JSON. Dans le cas d'un JSON :

bot_config

{
    "bot": {
        "token": 881hskdfASDA16641,
        "name": "Verter"
    },
    "service": "TG"
}

il suffit d'utiliser le chemin vers la clé souhaitée : {{ var.json.bot_config.bot.token }}.

Je vais dire un seul mot et montrer une capture d'écran sur les connexions. Ici, tout est élémentaire : sur la page Admin/Connections nous créons une connexion, en y ajoutant nos identifiants/mots de passe et des paramÚtres plus spécifiques. Voici comment :

Apache Airflow : simplifions l'ETL

Les mots de passe peuvent ĂȘtre chiffrĂ©s (plus soigneusement que dans l'option par dĂ©faut), ou le type de connexion peut ne pas ĂȘtre spĂ©cifiĂ© (comme je l'ai fait pour tg_main) — le fait est que la liste des types est intĂ©grĂ©e dans les modĂšles Airflow et ne peut ĂȘtre Ă©tendue sans toucher aux sources (si jamais je n'ai pas trouvĂ© l'information, merci de me corriger), mais rien ne nous empĂȘche d'obtenir les informations d'identification simplement par nom.

De plus, il est possible de crĂ©er plusieurs connexions avec le mĂȘme nom : dans ce cas, la mĂ©thode BaseHook.get_connection(), qui rĂ©cupĂšre les connexions par nom, renverra un alĂ©atoire parmi plusieurs homonymes (il serait plus logique de faire un Round Robin, mais laissons cela Ă  la discrĂ©tion des dĂ©veloppeurs d'Airflow).

Les Variables et les Connexions sont bien sĂ»r des outils pratiques, mais il est important de ne pas perdre l'Ă©quilibre : quelles parties de vos flux vous conservez effectivement dans le code, et quelles parties vous confiez Ă  Airflow. D'un cĂŽtĂ©, changer rapidement une valeur, par exemple une boĂźte de rĂ©ception, peut ĂȘtre pratique via l'UI. D'un autre cĂŽtĂ©, cela revient Ă  utiliser la souris, ce dont nous (je) voulions nous dĂ©barrasser.

Travailler avec les connexions est l'une des tùches hooks. En fait, les hooks Airflow sont des points de connexion à des services et bibliothÚques externes. Par exemple, JiraHook nous ouvrira un client pour interagir avec Jira (nous pouvons déplacer des tùches d'un cÎté à l'autre), et avec SambaHook nous pouvons pousser un fichier local vers un point smb.Nous sommes donc trÚs proches de voir comment fonctionne

Analyser un opérateur personnalisé

TelegramBotSendMessage commons/operators.py

Code avec l'opérateur proprement dit : avec son propre opérateur :

de typing import Union

from airflow.operators import BaseOperator

from commons.hooks import TelegramBotHook, TelegramBot

class TelegramBotSendMessage(BaseOperator):
    """Envoyer un message au chat_id en utilisant TelegramBotHook

    Exemple:
        >>> TelegramBotSendMessage(
        ...     task_id='telegram_fail', dag=dag,
        ...     tg_bot_conn_id='tg_bot_default',
        ...     chat_id='{{ var.value.all_the_young_dudes_chat }}',
        ...     message='{{ dag.dag_id }} a échoué :(',
        ...     trigger_rule=TriggerRule.ONE_FAILED)
    """
    template_fields = ['chat_id', 'message']

    def __init__(self,
                 chat_id: Union[int, str],
                 message: str,
                 tg_bot_conn_id: str = 'tg_bot_default',
                 *args, **kwargs):
        super().__init__(*args, **kwargs)

        self._hook = TelegramBotHook(tg_bot_conn_id)
        self.client: TelegramBot = self._hook.client
        self.chat_id = chat_id
        self.message = message

    def execute(self, context):
        print(f'Envoyer "{self.message}" au chat {self.chat_id}')
        self.client.send_message(chat_id=self.chat_id,
                                 message=self.message)

Ici, comme le reste dans Airflow, c'est trĂšs simple :

  • Nous avons hĂ©ritĂ© de BaseOperator, qui implĂ©mente beaucoup de choses spĂ©cifiques Ă  Airflow (jetez-y un Ɠil lorsque vous en avez le temps)
  • Nous avons dĂ©clarĂ© les champs template_fields, dans lesquels Jinja va chercher des macros pour le traitement.
  • Nous avons organisĂ© les bons arguments pour __init__(), en plaçant des valeurs par dĂ©faut lĂ  oĂč cela Ă©tait nĂ©cessaire.
  • Nous n'avons pas oubliĂ© d'initialiser le parent non plus.
  • Nous avons ouvert le hook correspondant TelegramBotHook, et obtenu l'objet client Ă  partir de celui-ci.
  • Nous avons redĂ©fini la mĂ©thode BaseOperator.execute(), que Airflow appellera lorsque le moment sera venu d'exĂ©cuter l'opĂ©rateur — c'est lĂ  que nous rĂ©alisons l'action principale, sans oublier de faire des journaux. (Nous loggons, d'ailleurs, directement dans stdout et stderr — Airflow va tout intercepter, bien l'enrober, le classer lĂ  oĂč il faut.)

Voyons ce que nous avons dans commons/hooks.py. La premiĂšre partie du fichier, avec le hook lui-mĂȘme :

from typing import Union

from airflow.hooks.base_hook import BaseHook
from requests_toolbelt.sessions import BaseUrlSession

class TelegramBotHook(BaseHook):
    """Hook API de Telegram Bot

    Remarque : ajoutez une connexion avec un type de Conn vide et n'oubliez pas
    de remplir Extra :

        {"bot_token": "YOuRAwEsomeBOtToKen"}
    """
    def __init__(self,
                 tg_bot_conn_id='tg_bot_default'):
        super().__init__(tg_bot_conn_id)

        self.tg_bot_conn_id = tg_bot_conn_id
        self.tg_bot_token = None
        self.client = None
        self.get_conn()

    def get_conn(self):
        extra = self.get_connection(self.tg_bot_conn_id).extra_dejson
        self.tg_bot_token = extra['bot_token']
        self.client = TelegramBot(self.tg_bot_token)
        return self.client

Je ne sais mĂȘme pas quoi expliquer ici, juste quelques points importants :

  • Nous hĂ©ritons, pensons aux arguments — dans la plupart des cas, il n'y en aura qu'un : conn_id;
  • Nous redĂ©finissons les mĂ©thodes standard : je me suis limitĂ© Ă  get_conn(), dans laquelle je rĂ©cupĂšre les paramĂštres de connexion par nom et j'extrais simplement la section extra (ce champ est pour JSON), dans lequel j'ai (selon mes propres instructions !) mis le token du bot Telegram : {"bot_token": "YOuRAwEsomeBOtToKen"}.
  • Je crĂ©e une instance de notre TelegramBot, lui donnant dĂ©jĂ  un token spĂ©cifique.

VoilĂ , c'est tout. On peut obtenir le client du hook grĂące Ă  TelegramBotHook().client ou TelegramBotHook().get_conn().

Et la deuxiĂšme partie du fichier, dans laquelle j'ai fait un micro-emballage pour l'API REST de Telegram, afin de ne pas avoir Ă  utiliser le mĂȘme python-telegram-bot pour une seule mĂ©thode sendMessage.

class TelegramBot:
    """Wrapper de l'API Telegram Bot

    Exemples:
        >>> TelegramBot('YOuRAwEsomeBOtToKen', '@myprettydebugchat').send_message('Salut, chéri')
        >>> TelegramBot('YOuRAwEsomeBOtToKen').send_message('Salut, chéri', chat_id=-1762374628374)
    """
    API_ENDPOINT = 'https://api.telegram.org/bot{}/'

    def __init__(self, tg_bot_token: str, chat_id: Union[int, str] = None):
        self._base_url = TelegramBot.API_ENDPOINT.format(tg_bot_token)
        self.session = BaseUrlSession(self._base_url)
        self.chat_id = chat_id

    def send_message(self, message: str, chat_id: Union[int, str] = None):
        method = 'sendMessage'

        payload = {'chat_id': chat_id or self.chat_id,
                   'text': message,
                   'parse_mode': 'MarkdownV2'}

        response = self.session.post(method, data=payload).json()
        if not response.get('ok'):
            raise TelegramBotException(response)

class TelegramBotException(Exception):
    def __init__(self, *args, **kwargs):
        super().__init__((args, kwargs))

Le bon chemin consiste Ă  rassembler tout cela : commons/operators.py, TelegramBotHook, TelegramBot — dans un plugin, le mettre dans un dĂ©pĂŽt public, et le rendre Open Source.

Pendant que nous étudions tout cela, nos mises à jour de rapports ont eu le temps de crouler sous le poids et de m'envoyer un message d'erreur dans le canal. Je vais vérifier ce qui ne va pas encore ...

Apache Airflow : simplifions l'ETL
Quelque chose s'est cassé dans notre DAG ! N'est-ce pas ce que nous attendions ? En effet !

Vas-tu remplir ça ?

Vous sentez que j'ai raté quelque chose ? Il semblait que je devais transférer des données de SQL Server vers Vertica, et voilà que je m'éloigne du sujet, quel drÎle !

Cet acte était intentionnel, je devais simplement vous expliquer une certaine terminologie. Maintenant, nous pouvons avancer.

Notre plan était le suivant :

  1. Faire le DAG
  2. Générer les tùches
  3. Regarder comme tout est beau
  4. Attribuer des numéros de session aux chargements
  5. Récupérer des données de SQL Server
  6. Mettre des données dans Vertica
  7. Rassembler des statistiques

Donc, pour démarrer tout cela, j'ai fait un petit ajout à notre docker-compose.yml:

docker-compose.db.yml

version: '3.4'

x-mssql-base: &mssql-base
  image: mcr.microsoft.com/mssql/server:2017-CU21-ubuntu-16.04
  restart: always
  environment:
    ACCEPT_EULA: Y
    MSSQL_PID: Express
    SA_PASSWORD: SayThanksToSatiaAt2020
    MSSQL_MEMORY_LIMIT_MB: 1024

services:
  dwh:
    image: jbfavre/vertica:9.2.0-7_ubuntu-16.04

  mssql_0:
    <<: *mssql-base

  mssql_1:
    <<: *mssql-base

  mssql_2:
    <<: *mssql-base

  mssql_init:
    image: mio101/py3-sql-db-client-base
    command: python3 .\/mssql_init.py
    depends_on:
      - mssql_0
      - mssql_1
      - mssql_2
    environment:
      SA_PASSWORD: SayThanksToSatiaAt2020
    volumes:
      - .\/mssql_init.py:\/mssql_init.py
      - .\/dags\/commons\/datasources.py:\/commons\/datasources.py

Ici, nous déployons :

  • Vertica en tant qu'hĂŽte dwh avec les paramĂštres par dĂ©faut les plus simples,
  • trois instances de SQL Server,
  • nous alimentons les bases avec quelques donnĂ©es rĂ©centes (surtout, ne regardez pas dans le mssql_init.py!)

Nous lançons tout cela avec une commande un peu plus complexe qu'auparavant :

$ docker-compose -f docker-compose.yml -f docker-compose.db.yml up --scale worker=3

Ce que notre gĂ©nĂ©rateur alĂ©atoire a produit peut ĂȘtre consultĂ© en utilisant la section Profilage des donnĂ©es / RequĂȘte Ad Hoc:

Apache Airflow : simplifions l'ETL
Surtout, ne le montrez pas aux analystes

Je ne vais pas m'attarder sur les sessions ETL c'est assez simple : nous créons une base de données, une table dans celle-ci, nous enveloppons le tout dans un gestionnaire de contexte, et maintenant faisons comme ceci :

with Session(task_name) as session:
    print('Chargement', session.id, 'commencé')

    # Charger le workflow
    ...

    session.successful = True
    session.loaded_rows = 15

session.py

from sys import stderr

class Session:
    """Session de flux de travail ETL

    Exemple:
        with Session(task_name) as session:
            print(session.id)
            session.successful = True
            session.loaded_rows = 15
            session.comment = 'Bien joué'
    """

    def __init__(self, connection, task_name):
        self.connection = connection
        self.connection.autocommit = True

        self._task_name = task_name
        self._id = None

        self.loaded_rows = None
        self.successful = None
        self.comment = None

    def __enter__(self):
        return self.open()

    def __exit__(self, exc_type, exc_val, exc_tb):
        if any(exc_type, exc_val, exc_tb):
            self.successful = False
            self.comment = f'{exc_type}: {exc_val}n{exc_tb}'
            print(exc_type, exc_val, exc_tb, file=stderr)
        self.close()

    def __repr__(self):
        return (f'')

    @property
    def task_name(self):
        return self._task_name

    @property
    def id(self):
        return self._id

    def _execute(self, query, *args):
        with self.connection.cursor() as cursor:
            cursor.execute(query, args)
            return cursor.fetchone()[0]

    def _create(self):
        query = """
            CREATE TABLE IF NOT EXISTS sessions (
                id          SERIAL       NOT NULL PRIMARY KEY,
                task_name   VARCHAR(200) NOT NULL,

                started     TIMESTAMPTZ  NOT NULL DEFAULT current_timestamp,
                finished    TIMESTAMPTZ           DEFAULT current_timestamp,
                successful  BOOL,

                loaded_rows INT,
                comment     VARCHAR(500)
            );
            """
        self._execute(query)

    def open(self):
        query = """
            INSERT INTO sessions (task_name, finished)
            VALUES (%s, NULL)
            RETURNING id;
            """
        self._id = self._execute(query, self.task_name)
        print(self, 'ouvert')
        return self

    def close(self):
        if not self._id:
            raise SessionClosedError('La session n'est pas ouverte')
        query = """
            UPDATE sessions
            SET
                finished    = DEFAULT,
                successful  = %s,
                loaded_rows = %s,
                comment     = %s
            WHERE
                id = %s
            RETURNING id;
            """
        self._execute(query, self.successful, self.loaded_rows,
                      self.comment, self.id)
        print(self, 'fermé',
              ', réussi: ', self.successful,
              ', Chargé: ', self.loaded_rows,
              ', commentaire:', self.comment)

class SessionError(Exception):
    pass

class SessionClosedError(SessionError):
    pass

Il est temps de récupérer nos données de nos quelque cent cinquante tables. Faisons-le avec quelques lignes trÚs simples :

source_conn = MsSqlHook(mssql_conn_id=src_conn_id, schema=src_schema).get_conn()

query = f"""
    SELECT 
        id, start_time, end_time, type, data
    FROM dbo.Orders
    WHERE
        CONVERT(DATE, start_time) = '{dt}'
    """

df = pd.read_sql_query(query, source_conn)
  1. Nous allons obtenir de Airflow pymssql-connexion
  2. Nous allons inclure une limite de date dans la requĂȘte - le templateur l'ajoutera.
  3. Nous alimentons notre requĂȘte pandas, qui nous extraira DataFrame — elle nous sera utile par la suite.

J'utilise le substitut {dt} au lieu du paramĂštre de la requĂȘte %s non pas parce que je suis un mĂ©chant Pinocchio, mais parce que pandas ne peut pas gĂ©rer pymssql et en propose au dernier params: Liste, bien qu'il le veuille vraiment tuple.
Veuillez également noter que le développeur pymssql a décidé de ne plus le supporter, et il est temps de passer à pyodbc.

Voyons ce qu'Airflow a introduit comme arguments Ă  nos fonctions :

Apache Airflow : simplifions l'ETL

S'il n'y a pas de données, il n'y a pas de raison de continuer. Mais il est également étrange de considérer que le chargement a réussi. Mais ce n'est pas une erreur. Que faire ? Voici la réponse :

if df.empty:
    raise AirflowSkipException('Aucune ligne Ă  charger')

AirflowSkipException Airflow dira qu'il n'y a pas vraiment d'erreur, et que nous passons la tùche. Dans l'interface, il n'y aura pas de carré vert ou rouge, mais de couleur rose.

Ajoutons à nos données quelques colonnes:

df['etl_source'] = src_schema
df['etl_id'] = session.id
df['hash_id'] = hash_pandas_object(df[['etl_source', 'id']])

À savoir :

  • BD d'oĂč nous avons rĂ©cupĂ©rĂ© les commandes,
  • L'identifiant de notre session de chargement (elle sera diffĂ©rente pour chaque tĂąche),
  • Le hachage de la source et de l'identifiant de la commande — pour que dans la base finale (oĂč tout sera fusionnĂ© dans une seule table) nous ayons un identifiant unique de commande.

Il reste une avant-derniĂšre Ă©tape : charger tout dans Vertica. Étonnamment, l'un des moyens les plus efficaces de le faire est via CSV !

# Export data to CSV buffer
buffer = StringIO()
df.to_csv(buffer,
          index=False, sep='|', na_rep='NUL', quoting=csv.QUOTE_MINIMAL,
          header=False, float_format='%.8f', doublequote=False, escapechar='\')
buffer.seek(0)

# Push CSV
target_conn = VerticaHook(vertica_conn_id=target_conn_id).get_conn()

copy_stmt = f"""
    COPY {target_table}({df.columns.to_list()}) 
    FROM STDIN 
    DELIMITER '|' 
    ENCLOSED '"' 
    ABORT ON ERROR 
    NULL 'NUL'
    """

cursor = target_conn.cursor()
cursor.copy(copy_stmt, buffer)
  1. Nous faisons un récepteur spécial StringIO.
  2. pandas va aimablement y stocker notre DataFrame sous forme CSV-lignes.
  3. Ouvrons une connexion Ă  notre chĂšre Vertica avec un hook.
  4. Et maintenant, à l'aide de copy() envoyons nos données directement dans Vertica !

Nous récupérerons depuis le driver combien de lignes ont été chargées, et dirons au manager de session que tout va bien :

session.loaded_rows = cursor.rowcount
session.successful = True

Et voilĂ .

En production, nous créons manuellement la table cible. Ici, je me suis permis un petit automatisme :

create_schema_query = f'CREATE SCHEMA IF NOT EXISTS {target_schema};'
create_table_query = f"""
    CREATE TABLE IF NOT EXISTS {target_schema}.{target_table} (
         id         INT,
         start_time TIMESTAMP,
         end_time   TIMESTAMP,
         type       INT,
         data       VARCHAR(32),
         etl_source VARCHAR(200),
         etl_id     INT,
         hash_id    INT PRIMARY KEY
     );"""

create_table = VerticaOperator(
    task_id='create_target',
    sql=[create_schema_query,
         create_table_query],
    vertica_conn_id=target_conn_id,
    task_concurrency=1,
    dag=dag)

À l'aide de VerticaOperator() je crĂ©e le schĂ©ma de BD et la table (s'ils n'existent pas encore, bien sĂ»r). L'essentiel est de bien Ă©tablir les dĂ©pendances :

for conn_id, schema in sql_server_ds:
    load = PythonOperator(
        task_id=schema,
        python_callable=workflow,
        op_kwargs={
            'src_conn_id': conn_id,
            'src_schema': schema,
            'dt': '{{ ds }}',
            'target_conn_id': target_conn_id,
            'target_table': f'{target_schema}.{target_table}'},
        dag=dag)

    create_table >> load

Bilan

— Eh bien, dit la petite souris, n'est-ce pas, maintenant
Es-tu sĂ»r que dans la forĂȘt, je suis la bĂȘte la plus terrifiante ?

Julia Donaldson, «Le Gruffalo»

Je pense que si mes collÚgues et moi organisions un concours : qui peut créer et lancer un processus ETL à partir de zéro le plus rapidement possible : eux avec leurs SSIS et leur souris et moi avec Airflow
 Et ensuite, nous comparerions la facilité de maintenance
 Ouf, je pense que vous conviendrez que je les surpasserai sur tous les fronts !

Cependant, pour parler un peu plus sĂ©rieusement, Apache Airflow — grĂące Ă  la description des processus sous forme de code — a rendu mon travail beaucoup plus pratique et agrĂ©able.

Sa capacitĂ© d'extension illimitĂ©e : tant en termes de plugins que de scalabilitĂ© — vous permet de l'appliquer pratiquement dans n'importe quel domaine : que ce soit dans le cycle complet de collecte, de prĂ©paration et de traitement des donnĂ©es, ou dans le lancement de fusĂ©es (vers Mars, bien sĂ»r).

Partie finale, informative et référentielle

Les écueils que nous avons évités pour vous

  • start_date. Oui, c'est dĂ©jĂ  un mĂšme local. À travers l'argument principal du DAG start_date tous passent. En rĂ©sumĂ©, si vous spĂ©cifiez dans start_date la date actuelle, et dans schedule_interval — un jour, alors le DAG ne se lancera pas avant demain.
    start_date = datetime(2020, 7, 7, 0, 1, 2)

    Et plus aucun problĂšme.

    C'est également lié à une autre erreur d'exécution : Task is missing the start_date parameter, qui indique le plus souvent que vous avez oublié de lier au DAG un opérateur.

  • Tout sur une seule machine. Oui, et les bases (d'Airflow lui-mĂȘme et de notre couvercle), ainsi que le serveur web, le planificateur et les workers. Et ça fonctionnait mĂȘme. Mais avec le temps, le nombre de tĂąches sur les services a augmentĂ©, et quand PostgreSQL a commencĂ© Ă  rĂ©pondre avec un index en 20 ms au lieu de 5 ms, nous l'avons pris et dĂ©placĂ©.
  • LocalExecutor. Oui, nous y sommes toujours, et nous avons dĂ©jĂ  atteint le bord du gouffre. Le LocalExecutor nous suffisait jusqu'Ă  prĂ©sent, mais maintenant il est temps de s'Ă©tendre d'au moins un worker, et il va falloir se bouger pour passer au CeleryExecutor. Et Ă©tant donnĂ© qu'on peut travailler avec lui mĂȘme sur une seule machine, rien n'empĂȘche l'utilisation de Celery mĂȘme sur un serveur qui « naturellement, n'ira jamais en production, c'est promis ! »
  • Ne pas utiliser les outils intĂ©grĂ©s:
    • Connections pour stocker les identifiants des services,
    • SLA Misses pour rĂ©agir aux tĂąches qui n'ont pas Ă©tĂ© exĂ©cutĂ©es Ă  temps,
    • XCom pour Ă©changer des mĂ©tadonnĂ©es (j'ai dit mĂ©tadonnĂ©es !) entre les tĂąches du DAG.
  • Abus de la messagerie. Que dire Ă  ce sujet ? Des alertes ont Ă©tĂ© configurĂ©es pour tous les doublons des tĂąches Ă©chouĂ©es. Maintenant, dans ma boĂźte Gmail de travail, j'ai plus de 90k emails d'Airflow, et l'interface web de la messagerie refuse de charger et de supprimer plus de 100 Ă  la fois.

Plus de piÚges cachés : Apache Airflow Pitfails

Outils pour encore plus d'automatisation

Pour nous permettre de travailler davantage avec notre tĂȘte plutĂŽt qu'avec nos mains, Airflow a prĂ©vu pour nous ceci :

  • API REST — il a toujours le statut Experimental, ce qui ne l'empĂȘche pas de fonctionner. Avec lui, il est possible non seulement d'obtenir des informations sur les DAGs et les tĂąches, mais aussi d'arrĂȘter/dĂ©marrer un DAG, de crĂ©er un DAG Run ou un pool.
  • CLI — de nombreux outils sont disponibles en ligne de commande, qui ne sont pas seulement difficiles Ă  utiliser via l'interface Web, mais qui sont mĂȘme totalement absents. Par exemple :
    • backfill nĂ©cessaire pour relancer des instances de tĂąches.
      Par exemple, des analystes arrivent et disent : « Eh bien, cher ami, il y a des anomalies dans les données du 1er au 13 janvier ! Réparez, réparez, réparez ! ». Et vous, vous répondez :
      airflow backfill -s '2020-01-01' -e '2020-01-13' orders
    • Maintenance de la base : initdb, resetdb, upgradedb, checkdb.
    • run, qui permet de lancer une seule instance de tĂąche, sans tenir compte de toutes les dĂ©pendances. De plus, il est possible de le lancer via LocalExecutor, mĂȘme si vous avez un cluster Celery.
    • En gros, c'est Ă  peu prĂšs ce que fait test, sauf qu'il n'Ă©crit rien dans la base.
    • connections permet de crĂ©er massivement des connexions depuis le shell.
  • Python API — un moyen assez hardcore d'interagir, qui est destinĂ© aux plugins, et non pas Ă  manipuler manuellement. Mais qui nous empĂȘche d'aller dans /home/airflow/dags, de lancer ipython et de commencer Ă  faire des folies ? On peut, par exemple, exporter toutes les connexions avec ce code :
    from airflow import settings
    from airflow.models import Connection
    
    fields = 'conn_id conn_type host port schema login password extra'.split()
    
    session = settings.Session()
    for conn in session.query(Connection).order_by(Connection.conn_id):
      d = {field: getattr(conn, field) for field in fields}
      print(conn.conn_id, '=', d)
  • Connexion Ă  la base de donnĂ©es des mĂ©tadonnĂ©es d'Airflow. Je ne recommande pas d'y Ă©crire, mais il est possible d'extraire les Ă©tats des tĂąches pour diffĂ©rentes mĂ©triques spĂ©cifiques beaucoup plus rapidement et facilement que via n'importe quel des API.

    Disons que tous nos travaux ne sont pas idempotents et peuvent parfois échouer, et c'est normal. Mais plusieurs échecs, c'est déjà suspect, et il faudrait vérifier.

    Attention, SQL !

    AVEZ derniÚres_exécutions COMME (
    SELECT
        task_id,
        dag_id,
        execution_date,
        state,
            row_number()
            OVER (
                PARTITION PAR task_id, dag_id
                ORDER BY execution_date DESC) AS rn
    FROM public.task_instance
    WHERE
        execution_date > now() - INTERVAL '2' JOURS
    ),
    échoué COMME (
        SELECT
            task_id,
            dag_id,
            execution_date,
            state,
            CASE WHEN rn = row_number() OVER (
                PARTITION PAR task_id, dag_id
                ORDER BY execution_date DESC)
                     THEN TRUE END AS last_fail_seq
        FROM derniÚres_exécutions
        WHERE
            state IN ('failed', 'up_for_retry')
    )
    SELECT
        task_id,
        dag_id,
        count(last_fail_seq)                       AS unsuccessful,
        count(CASE WHEN last_fail_seq
            AND state = 'failed' THEN 1 END)       AS failed,
        count(CASE WHEN last_fail_seq
            AND state = 'up_for_retry' THEN 1 END) AS up_for_retry
    FROM échoué
    GROUPE PAR
        task_id,
        dag_id
    AYANT
        count(last_fail_seq) > 0

Liens

Et bien sûr, les dix premiers liens dans les résultats de Google contiennent le contenu du dossier Airflow de mes favoris.

Et les liens mentionnés dans l'article :

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