Airflow è uno strumento per sviluppare e mantenere in modo facile e veloce i processi batch di elaborazione dati

Airflow è uno strumento per sviluppare e mantenere in modo facile e veloce i processi batch di elaborazione dati

Ciao, Habr! In questo articolo voglio parlarvi di uno strumento straordinario per lo sviluppo di processi batch per l'elaborazione dei dati, ad esempio, nell'infrastruttura di un DWH aziendale o del vostro DataLake. Parleremo di Apache Airflow (di seguito Airflow). È ingiustamente trascurato su Habr e nella parte principale cercherò di convincervi che almeno Airflow merita attenzione quando si sceglie un pianificatore per i vostri processi ETL/ELT.

In precedenza ho scritto una serie di articoli sul tema DWH quando lavoravo in Tinkoff Bank. Ora sono diventato parte del team di Mail.Ru Group e mi occupo dello sviluppo della piattaforma per l'analisi dei dati nella direzione dei giochi. In effetti, man mano che ci saranno notizie e soluzioni interessanti, io e il mio team parleremo qui della nostra piattaforma per l'analisi dei dati.

Prologo

Quindi, cominciamo. Cos'è Airflow? È una libreria (o un insieme di librerie) per lo sviluppo, la pianificazione e il monitoraggio dei flussi di lavoro. La caratteristica principale di Airflow è che per descrivere (sviluppare) i processi si utilizza codice in linguaggio Python. Da ciò derivano numerosi vantaggi per organizzare il vostro progetto e lo sviluppo: essenzialmente, il vostro progetto ETL, ad esempio, è semplicemente un progetto Python, e potete organizzarne la struttura come preferite, considerando le caratteristiche dell'infrastruttura, le dimensioni del team e altri requisiti. Strumentalmente è tutto semplice. Utilizzate, ad esempio, PyCharm + Git. È fantastico e molto comodo!

Ora esaminiamo le principali entità di Airflow. Comprendendo la loro essenza e il loro scopo, potrete organizzare al meglio l'architettura dei processi. Probabilmente, l'entità principale è il Directed Acyclic Graph (di seguito DAG).

DAG

Il DAG è una certa unione semantica delle vostre attività che desiderate eseguire in un'ordine rigorosamente definito secondo un programma specifico. Airflow offre un'interfaccia web comoda per lavorare con i DAG e le altre entità:

Airflow è uno strumento per sviluppare e mantenere in modo facile e veloce i processi batch di elaborazione dati

Un DAG può apparire in questo modo:

Airflow è uno strumento per sviluppare e mantenere in modo facile e veloce i processi batch di elaborazione dati

Lo sviluppatore, progettando un DAG, fissa un insieme di operatori su cui saranno basate le attività all'interno del DAG. Qui arriviamo a un'altra entità importante: l'Operatore di Airflow.

Operatori

L'Operatore è un'entità sulla base della quale vengono creati gli istanti di lavoro, in cui si descrive cosa accadrà durante l'esecuzione dell'istanza di lavoro. Le versioni di Airflow su GitHub contengono già un insieme di operatori pronti all'uso. Esempi:

  • BashOperator — operatore per l'esecuzione di comandi bash.
  • PythonOperator — operatore per chiamare codice Python.
  • EmailOperator — operatore per inviare email.
  • HTTPOperator — operatore per l'interazione con le richieste http.
  • SqlOperator — operatore per eseguire codice SQL.
  • Sensor — operatore che attende un evento (come un tempo specifico, l'arrivo di un file richiesto, una riga in un database, una risposta da un'API, ecc.).

Ci sono operatori più specifici: DockerOperator, HiveOperator, S3FileTransferOperator, PrestoToMysqlOperator, SlackOperator.

È anche possibile sviluppare operatori in base alle proprie esigenze e utilizzarli nel progetto. Ad esempio, abbiamo creato MongoDBToHiveViaHdfsTransfer, un operatore per l'esportazione di documenti da MongoDB a Hive, e diversi operatori per lavorare con ClickHouse: CHLoadFromHiveOperator e CHTableLoaderOperator. Fondamentalmente, quando in un progetto appare un codice frequentemente utilizzato costruito su operatori di base, è possibile considerare di raccoglierlo in un nuovo operatore. Questo semplificherà lo sviluppo futuro e arricchirà la tua libreria di operatori nel progetto.

Dopo, tutte queste istanze di attività devono essere eseguite, e ora parleremo del pianificatore.

Pianificatore

Il pianificatore di attività in Airflow è basato su Celery. Celery è una libreria Python che consente di organizzare una coda e l'esecuzione asincrona e distribuita delle attività. Dal lato di Airflow tutte le attività sono suddivise in pool. I pool vengono creati manualmente. Di solito, il loro obiettivo è limitare il carico di lavoro sull'origine o categorizzare le attività all'interno del DWH. I pool possono essere gestiti tramite un'interfaccia web:

Airflow è uno strumento per sviluppare e mantenere in modo facile e veloce i processi batch di elaborazione dati

Ogni pool ha un limite sul numero di slot. Durante la creazione di un DAG, viene assegnato un pool:

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

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

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

Il pool assegnato a livello di DAG può essere sovrascritto a livello di attività.
La pianificazione di tutte le attività in Airflow è gestita da un processo separato — Scheduler. Fondamentalmente, lo Scheduler si occupa di tutta la meccanica di assegnazione delle attività all'esecuzione. Un'attività, prima di essere eseguita, passa attraverso diverse fasi:

  1. Nel DAG sono state completate le attività precedenti, la nuova può essere messa in coda.
  2. La coda è ordinata in base alla priorità delle attività (le priorità possono essere gestite anche) e, se c'è uno slot libero nel pool, l'attività può essere presa in carico.
  3. Se c'è un worker celery libero, l'attività viene inviata a lui; inizia il lavoro che hai programmato nell'attività, utilizzando un determinato operatore.

È abbastanza semplice.

Lo Scheduler funziona su molti DAG e su tutte le attività all'interno dei DAG.

Affinché lo Scheduler inizi a lavorare con il DAG, è necessario impostare un programma per il DAG:

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

Ci sono alcuni preset già pronti: @once, @hourly, @daily, @weekly, @monthly, @yearly.

È possibile utilizzare anche le espressioni cron:

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

Data di Esecuzione

Per capire come funziona Airflow, è importante sapere che cos'è la Data di Esecuzione per il DAG. In Airflow, il DAG ha una dimensione Data di Esecuzione, cioè, in base al programma di lavoro del DAG, vengono creati esemplari delle attività per ogni Data di Esecuzione. E per ogni Data di Esecuzione, le attività possono essere eseguite nuovamente — o, ad esempio, il DAG può lavorare simultaneamente in più Date di Esecuzione. Questo è chiaramente mostrato qui:

Airflow è uno strumento per sviluppare e mantenere in modo facile e veloce i processi batch di elaborazione dati

Sfortunatamente (o forse anche fortunate: dipende dalla situazione), se viene modificata l'implementazione dell'attività nel DAG, l'esecuzione nelle Date di Esecuzione precedenti avverrà già tenendo conto delle correzioni. Questo è positivo se è necessario rielaborare i dati nei periodi passati con un nuovo algoritmo, ma è negativo perché si perde la riproducibilità del risultato (certo, nessuno ti impedisce di recuperare la versione necessaria del sorgente da Git e calcolare all'occorrenza ciò che è necessario, esattamente come serve).

Generazione delle attività

L'implementazione del DAG è codice Python, quindi abbiamo un modo molto comodo per ridurre il volume del codice quando lavoriamo, ad esempio, con fonti shardate. Supponiamo che tu abbia come fonte tre shard MySQL, devi accedere a ciascuno e recuperare alcuni dati. E questo in modo indipendente e parallelo. Il codice in Python nel DAG potrebbe apparire così:

connection_list = lv.get('connection_list')

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

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

Il DAG risulta così:

Airflow è uno strumento per sviluppare e mantenere in modo facile e veloce i processi batch di elaborazione dati

È possibile aggiungere o rimuovere uno shard, semplicemente regolando le impostazioni e aggiornando il DAG. Comodo!

È possibile utilizzare anche generazioni di codice più complesse, ad esempio lavorare con fonti come database o descrivere la struttura delle tabelle, l'algoritmo di lavoro con la tabella e, tenendo conto delle peculiarità dell'infrastruttura DWH, generare il processo di caricamento di N tabelle nel vostro magazzino. Oppure, ad esempio, lavorare con un'API che non supporta il lavoro con parametri sotto forma di lista, è possibile generare N task nel DAG da questa lista, limitando la concorrenza nelle richieste all'API tramite un pool e raccogliendo i dati necessari dall'API. Flessibile!

Repository

In Airflow c'è un proprio repository backend, un database (può essere MySQL o Postgres, noi usiamo Postgres), in cui vengono memorizzati gli stati dei task, dei DAG, le impostazioni delle connessioni, le variabili globali, ecc. Qui vorrei sottolineare che il repository in Airflow è molto semplice (circa 20 tabelle) e comodo, se volete costruire un qualsiasi processo su di esso. Ricordo le 100500 tabelle nel repository di Informatica, che richiedevano molto tempo per essere comprese prima di capire come costruire una query.

Monitoraggio

Data la semplicità del repository, potete costruire un processo di monitoraggio dei task che risponda alle vostre esigenze. Noi utilizziamo un notebook in Zeppelin, dove controlliamo lo stato dei task:

Airflow è uno strumento per sviluppare e mantenere in modo facile e veloce i processi batch di elaborazione dati

Questo può essere anche l'interfaccia web di Airflow:

Airflow è uno strumento per sviluppare e mantenere in modo facile e veloce i processi batch di elaborazione dati

Il codice di Airflow è aperto, quindi abbiamo aggiunto allerta tramite Telegram. Ogni istanza di task attiva, in caso di errore, invia messaggi nel gruppo Telegram, dove è presente tutto il team di sviluppo e supporto.

Riceviamo risposte tempestive tramite Telegram (se necessario), mentre attraverso Zeppelin otteniamo una visione complessiva sui task in Airflow.

Totale

Airflow è principalmente open source e non aspettatevi miracoli. Siate pronti a investire tempo e sforzi per costruire una soluzione funzionante. L'obiettivo è raggiungibile, credetemi, ne vale la pena. La velocità di sviluppo, la flessibilità, la semplicità di aggiunta di nuovi processi: vi piacerà. Naturalmente, è necessario prestare molta attenzione all'organizzazione del progetto e alla stabilità del funzionamento di Airflow: non ci sono miracoli.

Attualmente abbiamo Airflow che esegue quotidianamente circa 6.500 task. Per natura, sono abbastanza diversi. Ci sono compiti di caricamento dei dati nel DWH principale da molteplici fonti diverse e molto specifiche, ci sono compiti di calcolo dei dati all'interno del DWH principale, ci sono compiti di pubblicazione dei dati in un DWH veloce, ci sono molte, molte diverse mansioni — e Airflow le gestisce giorno dopo giorno. Se parliamo in numeri, sono 2.3 mila compiti ELT di diversa complessità all'interno del DWH (Hadoop), circa 2.5 centinaia di database fonti, questo è un team di 4 sviluppatori ETL, che si dividono tra il processamento ETL dei dati nel DWH e il processamento ELT dei dati all'interno del DWH e naturalmente anche un amministratore, che si occupa dell'infrastruttura del servizio.

Piani futuri

Il numero dei processi cresce inevitabilmente, e la cosa principale su cui ci concentreremo riguardo all'infrastruttura di Airflow è la scalabilità. Vogliamo costruire un cluster di Airflow, dedicare un paio di nodi ai worker di Celery e creare una testa duplicata con processi di pianificazione delle attività e un repository.

Epilogo

Questo, ovviamente, non è tutto ciò che vorrei raccontare su Airflow, ma ho cercato di coprire i punti principali. L'appetito viene mangiando, prova — e ti piacerà 🙂

Fonte: habr.com

Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server 🔥 Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server | ProHoster