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.

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

- 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 :
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
. /dagsnous 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. .
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
- brokerRemarques :
- Dans la construction du compose, je me suis largement basĂ© sur l'image bien connue â Ă 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=3AprĂšs que tout s'est mis en route, vous pouvez consulter les interfaces web :
- Airflow :
- Flower :
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.
HttpSensorpeut interroger un endpoint spĂ©cifiĂ©, et lorsque la bonne rĂ©ponse est reçue, lancer le transfertGoogleCloudStorageToS3Operator. 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.
- action, comme notre cher
- 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
LocalExecutorou sur un nĆud distant dans le cas deCeleryExecutor), 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 desList[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 dansglobals(), 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 unecron-chaĂźne0 0 0/6 ? * * *, pour les moins durs â une expression comme@daily);
- quel est son nom
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 videworkflow(). 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 flagprovide_contextinjectera Ă 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.

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 :

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 :

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 dansgitnotre Gitlab, et Gitlab CI dĂ©ploie les mises Ă jour sur les machines lors d'un merge dansmaster.
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 :

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

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

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

Charger ce qui n'est pas encore chargé
Ainsi, toutes les tùches ont été effectuées, nous pouvons évacuer les blessés.

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 :

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.

Ă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

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

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

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 }}'â champtonous 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 dansAdmin/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'â argumentsconn_idreçoivent les identifiants de connexions que nous crĂ©ons dansAdmin/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 plusieursVerticaOperator(regardant la mĂȘme table);report_update >> [email, tg]â toutVerticaOperatorse rejoindra pour l'envoi de l'email et du message, comme ceci :

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 :

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 :

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

Une liste complÚte des macros intégrées pour la derniÚre version disponible est disponible ici :
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 :

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 :

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 dansstdoutetstderrâ 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.clientJe 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 sectionextra(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 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 ...

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 :
- Faire le DAG
- Générer les tùches
- Regarder comme tout est beau
- Attribuer des numéros de session aux chargements
- Récupérer des données de SQL Server
- Mettre des données dans Vertica
- 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.pyIci, nous déployons :
- Vertica en tant qu'hĂŽte
dwhavec 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=3Ce que notre gĂ©nĂ©rateur alĂ©atoire a produit peut ĂȘtre consultĂ© en utilisant la section Profilage des donnĂ©es / RequĂȘte Ad Hoc:

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 = 15session.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):
passIl 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)- Nous allons obtenir de Airflow
pymssql-connexion - Nous allons inclure une limite de date dans la requĂȘte - le templateur l'ajoutera.
- Nous alimentons notre requĂȘte
pandas, qui nous extrairaDataFrameâ elle nous sera utile par la suite.
J'utilise le substitut
{dt}au lieu du paramĂštre de la requĂȘte%snon pas parce que je suis un mĂ©chant Pinocchio, mais parce quepandasne peut pas gĂ©rerpymssqlet en propose au dernierparams: Liste, bien qu'il le veuille vraimenttuple.
Veuillez Ă©galement noter que le dĂ©veloppeurpymssqla dĂ©cidĂ© de ne plus le supporter, et il est temps de passer Ăpyodbc.
Voyons ce qu'Airflow a introduit comme arguments Ă nos fonctions :

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)- Nous faisons un récepteur spécial
StringIO. pandasva aimablement y stocker notreDataFramesous formeCSV-lignes.- Ouvrons une connexion Ă notre chĂšre Vertica avec un hook.
- 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 = TrueEt 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 >> loadBilan
â 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 DAGstart_datetous passent. En rĂ©sumĂ©, si vous spĂ©cifiez dansstart_datela date actuelle, et dansschedule_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 :
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 :
- â 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.
- â 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 :
backfillné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 viaLocalExecutor, 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. connectionspermet de créer massivement des connexions depuis le shell.
- â 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 lanceripythonet 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.
- â bien sĂ»r, il faut commencer par la documentation officielle, mais qui lit les instructions ?
- â au moins lisez les recommandations des crĂ©ateurs.
- â les bases : l'interface utilisateur en images
- â les bases sont bien expliquĂ©es, au cas oĂč (par hasard !) vous n'auriez pas compris chez moi.
- â un guide rapide pour configurer un cluster Airflow.
- â un article presque aussi intĂ©ressant, juste un peu plus de formalisme et moins d'exemples.
- â sur le travail conjoint avec Celery.
- â sur l'idempotence des tĂąches, le chargement par ID au lieu de la date, les transformations, la structure des fichiers et d'autres choses intĂ©ressantes.
- â dĂ©pendances des tĂąches et Trigger Rule, que j'ai mentionnĂ©es briĂšvement.
- â comment surmonter certains « fonctionne comme prĂ©vu » du planificateur, charger les donnĂ©es manquantes et prioriser les tĂąches.
- â requĂȘtes SQL utiles pour les mĂ©tadonnĂ©es Airflow.
- â il y a une section utile sur la crĂ©ation de capteurs personnalisĂ©s.
- â une note courte et intĂ©ressante sur la construction de l'infrastructure sur AWS pour Data Science.
- â erreurs frĂ©quentes (lorsque quelqu'un ne lit toujours pas les instructions).
- â souriez, comment les gens bricolent le stockage des mots de passe, alors qu'il suffit d'utiliser les Connections.
- â passage implicite du DAG, passage de contexte dans une fonction, encore sur les dĂ©pendances, et aussi sur le saut des exĂ©cutions de tĂąches.
- â sur l'utilisation de
arguments par dĂ©fautetparamsdans les modĂšles, ainsi que sur les variables et les connexions. - â un rĂ©cit sur la prĂ©paration du planificateur pour Airflow 2.0.
- â un article un peu obsolĂšte sur le dĂ©ploiement de notre cluster dans
docker-compose. - â des tĂąches dynamiques Ă l'aide de modĂšles et de passage de contexte.
- â notifications standard et personnalisĂ©es par e-mail et Slack.
- â embranchements de tĂąches, macros et XCom.
Et les liens mentionnés dans l'article :
- â placeholders disponibles pour utilisation dans les modĂšles.
- â erreurs frĂ©quentes lors de la crĂ©ation de DAG.
- â
docker-composepour les expĂ©riences, le dĂ©bogage et plus encore. - â Wrapper Python pour l'API REST de Telegram.
Source : habr.com




