Cómo hacer un disparador DAG en Airflow usando la API experimental

En la elaboración de nuestros programas educativos, periódicamente nos encontramos con dificultades en cuanto al trabajo con algunas herramientas. Y en el momento en que nos encontramos con ellos, no siempre hay suficiente documentación y artículos que ayuden a hacer frente a este problema.

Así fue, por ejemplo, en 2015, y usamos el clúster Hadoop con Spark para 35 usuarios simultáneos en el programa Big Data Specialist. No estaba claro cómo prepararlo para tal caso de usuario usando YARN. Como resultado, habiendo descubierto y caminando el camino por su cuenta, no publicar en Habré y también realizó Reunión Spark de Moscú.

Prehistoria

Esta vez hablaremos de un programa diferente: Data Engineer. Sobre él, nuestros participantes construyen dos tipos de arquitectura: lambda y kappa. Y en la arquitectura lamdba, Airflow se usa como parte del procesamiento por lotes para transferir registros de HDFS a ClickHouse.

Todo está bien en general. Que construyan sus propios oleoductos. Sin embargo, hay un pero: todos nuestros programas son tecnológicamente avanzados desde el punto de vista del propio proceso de aprendizaje. Para verificar el laboratorio utilizamos verificadores automáticos: el participante debe ir a su cuenta personal, hacer clic en el botón "Verificar" y después de un tiempo ve algún tipo de retroalimentación ampliada sobre lo que hizo. Y es en este momento cuando comenzamos a abordar nuestro problema.

La verificación de este laboratorio se organiza de la siguiente manera: enviamos un paquete de datos de control al Kafka del participante, luego Gobblin transfiere este paquete de datos a HDFS, luego Airflow toma este paquete de datos y lo coloca en ClickHouse. El truco es que Airflow no tiene que hacer esto en tiempo real, lo hace según lo programado: una vez cada 15 minutos, toma un montón de archivos y los sube.

Resulta que necesitamos activar de alguna manera su DAG por nuestra cuenta a petición nuestra mientras el verificador se está ejecutando aquí y ahora. Buscando en Google, descubrimos que para versiones posteriores de Airflow existe un llamado API experimental. palabra experimental, por supuesto, suena aterrador, pero qué hacer ... De repente despega.

A continuación, describiremos todo el proceso: desde la instalación de Airflow hasta la creación de una solicitud POST que activa un DAG utilizando la API experimental. Trabajaremos con Ubuntu 16.04.

1. Instalación de flujo de aire

Comprobemos que tenemos Python 3 y virtualenv.

$ python3 --version
Python 3.6.6
$ virtualenv --version
15.2.0

Si falta uno de estos, instálelo.

Ahora vamos a crear un directorio en el que seguiremos trabajando con Airflow.

$ mkdir <your name of directory>
$ cd /path/to/your/new/directory
$ virtualenv -p which python3 venv
$ source venv/bin/activate
(venv) $

Instale el flujo de aire:

(venv) $ pip install airflow

Versión en la que trabajamos: 1.10.

Ahora necesitamos crear un directorio. airflow_home, donde se ubicarán los archivos DAG y los complementos de Airflow. Después de crear el directorio, configure la variable de entorno AIRFLOW_HOME.

(venv) $ cd /path/to/my/airflow/workspace
(venv) $ mkdir airflow_home
(venv) $ export AIRFLOW_HOME=<path to airflow_home>

El siguiente paso es ejecutar el comando que creará e inicializará la base de datos de flujo de datos en SQLite:

(venv) $ airflow initdb

La base de datos se creará en airflow.db defecto.

Compruebe si Airflow está instalado:

$ airflow version
[2018-11-26 19:38:19,607] {__init__.py:57} INFO - Using executor SequentialExecutor
[2018-11-26 19:38:19,745] {driver.py:123} INFO - Generating grammar tables from /usr/lib/python3.6/lib2to3/Grammar.txt
[2018-11-26 19:38:19,771] {driver.py:123} INFO - Generating grammar tables from /usr/lib/python3.6/lib2to3/PatternGrammar.txt
  ____________       _____________
 ____    |__( )_________  __/__  /________      __
____  /| |_  /__  ___/_  /_ __  /_  __ _ | /| / /
___  ___ |  / _  /   _  __/ _  / / /_/ /_ |/ |/ /
 _/_/  |_/_/  /_/    /_/    /_/  ____/____/|__/
   v1.10.0

Si el comando funcionó, Airflow creó su propio archivo de configuración airflow.cfg в AIRFLOW_HOME:

$ tree
.
├── airflow.cfg
└── unittests.cfg

Airflow tiene una interfaz web. Se puede iniciar ejecutando el comando:

(venv) $ airflow webserver --port 8081

Ahora puede acceder a la interfaz web en un navegador en el puerto 8081 en el host donde se estaba ejecutando Airflow, así: <hostname:8081>.

2. Trabajando con la API Experimental

En este Airflow está configurado y listo para funcionar. Sin embargo, también necesitamos ejecutar la API experimental. Nuestras damas están escritas en Python, por lo que todas las solicitudes se realizarán utilizando la biblioteca. requests.

En realidad, la API ya está funcionando para solicitudes simples. Por ejemplo, dicha solicitud le permite probar su trabajo:

>>> import requests
>>> host = <your hostname>
>>> airflow_port = 8081 #в нашем случае такой, а по дефолту 8080
>>> requests.get('http://{}:{}/{}'.format(host, airflow_port, 'api/experimental/test').text
'OK'

Si recibió un mensaje de este tipo en respuesta, significa que todo está funcionando.

Sin embargo, cuando queremos activar un DAG, nos encontramos con el hecho de que este tipo de solicitud no se puede realizar sin autenticación.

Para hacer esto, deberá realizar una serie de acciones.

Primero, debe agregar esto a la configuración:

[api]
auth_backend = airflow.contrib.auth.backends.password_auth

Luego, debe crear su usuario con derechos de administrador:

>>> import airflow
>>> from airflow import models, settings
>>> from airflow.contrib.auth.backends.password_auth import PasswordUser
>>> user = PasswordUser(models.Admin())
>>> user.username = 'new_user_name'
>>> user.password = 'set_the_password'
>>> session = settings.Session()
>>> session.add(user)
>>> session.commit()
>>> session.close()
>>> exit()

A continuación, debe crear un usuario con derechos normales que pueda activar un DAG.

>>> import airflow
>>> from airflow import models, settings
>>> from airflow.contrib.auth.backends.password_auth import PasswordUser
>>> user = PasswordUser(models.User())
>>> user.username = 'newprolab'
>>> user.password = 'Newprolab2019!'
>>> session = settings.Session()
>>> session.add(user)
>>> session.commit()
>>> session.close()
>>> exit()

Ahora todo está listo.

3. Lanzar una solicitud POST

La solicitud POST en sí se verá así:

>>> dag_id = newprolab
>>> url = 'http://{}:{}/{}/{}/{}'.format(host, airflow_port, 'api/experimental/dags', dag_id, 'dag_runs')
>>> data = {"conf":"{"key":"value"}"}
>>> headers = {'Content-type': 'application/json'}
>>> auth = ('newprolab', 'Newprolab2019!')
>>> uri = requests.post(url, data=json.dumps(data), headers=headers, auth=auth)
>>> uri.text
'{n  "message": "Created <DagRun newprolab @ 2019-03-27 10:24:25+00:00: manual__2019-03-27T10:24:25+00:00, externally triggered: True>"n}n'

Solicitud procesada con éxito.

En consecuencia, le damos al DAG algo de tiempo para que procese y haga una solicitud a la tabla de ClickHouse, tratando de capturar el paquete de datos de control.

Verificación completada.

Fuente: habr.com

Compre alojamiento confiable para sitios con protección DDoS, servidores VPS VDS 🔥 Compra alojamiento web fiable con protección DDoS, servidores VPS VDS | ProHoster