Apache Airflow: teeme ETL lihtsamaks

Tere, mina olen Dmitri Logvinenko — andengineer grupi firmade «Vezot» analĂŒĂŒtika osakonnas.

Kojutame teile suurepĂ€rasest tööriistast ETL-protsesside arendamiseks — Apache Airflow. Kuid Airflow on nii universaalne ja mitmekesine, et tasub sellele tĂ€helepanu pöörata ka siis, kui te ei tegele andmevoogudega, vaid vajate aeg-ajalt mĂ”ne protsessi kĂ€ivitamist ja nende tĂ€itmise jĂ€lgimist.

Jah, ma ei rÀÀgi ainult, vaid nÀitan ka: programmi sisaldab palju koodi, ekraanipilte ja soovitusi.

Apache Airflow: teeme ETL lihtsamaks
Mida tavaliselt nÀed, kui googeldad sÔna Airflow / Wikimedia Commons

Sisukord

Sissejuhatus

Apache Airflow — ta on nagu Django:

  • kirjutatud Pythonis,
  • seal on suurepĂ€rane adminpaneel,
  • piiramatu laiendatavusega,

— ainult parim, ja see on loodud hoopis muude eesmĂ€rkide jaoks, nimelt (nagu enne katset):

  • ĂŒlesannete kĂ€ivitamine ja jĂ€lgimine piiramatul arvul masinatel (nii palju kui teie sĂŒdametunnistus ja Celery/Kubernetes lubavad)
  • dĂŒnaamilise töövoo genereerimise abil, kasutades vĂ€ga lihtsat Python-koodi, mis on lihtne kirjutada ja mĂ”ista
  • ning vĂ”imalusega siduda omavahel mis tahes andmebaasid ja API-d valmiskomponentide ning isetehtud pluginatega (mida on ÀÀrmiselt lihtne teha).

Kasutame Apache Airflow'i jÀrgmiselt:

  • kogume andmeid erinevatest allikatest (mitmed SQL Serveri ja PostgreSQL instance, erinevad API-d rakenduste mÔÔdikute jaoks, isegi 1C) DWH-sse ja ODS-sse (meie puhul Vertica ja Clickhouse).
  • nagu arenenud cron, mis kĂ€ivitab andmete konsolideerimise protsesse ODS-is ja jĂ€lgib ka nende haldamist.

Kuni hiljuti katab meie vajadusi ĂŒks vĂ€ike server 32 tuuma ja 50 GB mĂ€luga. Airflow'is töötab sel juhul:

  • ĂŒle 200 dag'i (töötavad workflows, kuhu oleme ĂŒlesandeid tĂ€itnud),
  • igaĂŒhes keskmiselt 70 ĂŒlesannet,
  • killedakse seda (jĂ€llegi keskmiselt) ĂŒks kord tunnis..

Ja sellest, kuidas me laienesime, kirjutan alla, aga nĂŒĂŒd mÀÀratleme ĂŒber-ĂŒlesande, mida me lahendame:

Kolmes SQL Serverit, igaĂŒhel 50 andmebaasi – ĂŒhe projekti instantsid, seega on neil struktuur peaaegu sama (peaaegu kĂ”ikjal, muha-haha), mis tĂ€hendab, et igas on tabel Orders (Ă”nneks saab sellise nimega tabelit panna igasse Ă€risse). Me vĂ”tame andmed, lisades teenindusvĂ€ljad (allikaserver, allikandmebaas, ETL-ĂŒlesande identifikaator) ning ĂŒritame nad sĂŒĂŒdistusteta visata, ĂŒtleme, Verticasse.

LĂ€hme!

Peamine, praktiline osa (ja natuke teooriat)

Miks see meile (ja teile) vajalik on

Kui puud olid suured ja mina olin lihtne SQL-hooldaja ĂŒhes Venemaa jaeketis, heitsime ETL-protsesside ehk andmevoogude kallale kahe kĂ€esolevate vahenditega:

  • Informatica Power Center — ÀÀrmiselt keeruline sĂŒsteem, ÀÀrmiselt tootlik, oma riistvara ja versioneerimisega. Kasutasin vĂ”ib-olla 1% selle vĂ”imalustest. Miks? No, esmalt, see liides on kuskil nullindatest vaimselt tapnud meid. Teiseks, see vidin on loodud ÀÀrmiselt keerukate protsesside, intensiivse komponentide taaskasutamise ja teiste vĂ€ga-oluliste ettevĂ”tte-nippide jaoks. Me jĂ€tame vahele, et see maksab sama palju kui Airbus A380 tiib/ aasta.

    Olge ettevaatlik, ekraanipilt vÔib inimestele, kes on alla 30, natuke haiget teha.

    Apache Airflow: teeme ETL lihtsamaks

  • SQL Server Integration Server — seda tööriista oleme kasutanud oma siseprojektide voogudes. Tegelikult kasutame me juba SQL Serverit ja selle ETL tööriistade mittekasutamine oleks olnud ebaaus. KĂ”ik on selles hĂ€sti: liides on ilus, ja töötlemisaruanded... Aga me ei armasta tarkvaratooteid selle pĂ€rast, oh ei. Versioonide haldamine dtsx (mis on XML, millel on salvestamisel segased sĂ”lmed) on meil vĂ”imalik, aga mis sellest kasu? Kuidas luua ĂŒlesannete pakett, mis tĂ”mbab sada tabelit ĂŒhest serverist teise? Mis sada, kahekĂŒmne tabeliga kaotab teil sĂ”rm, mis klĂ”psab rottnuppu. Kuid see nĂ€eb kindlasti stiilsem vĂ€lja:

    Apache Airflow: teeme ETL lihtsamaks

Oleme absoluutsetelt otsinud vÀljapÀÀse. Asi oli peaaegu peaaegu jÔudnud kohandatud SSIS-pakettide generaatorini...

... ja siis leidis mind uus töö. Ja seal tabas mind Apache Airflow.

Kui ma sain teada, et ETL-protsesside kirjeldused on lihtsalt Python'i kood, ei osanud ma rÔÔmust tantsida. Nii andmevood lĂ€bisid versioonihaldust ja diffimist, ning tabelite koondamine ĂŒhise struktuuriga sadadest andmebaasidest ĂŒhte sihtkohta muutus Python'i koodi kettaks 13” ekraanil.

Klastri kokkupanemine

Ärge muretsege siin lasteaiatuks, rÀÀgime hoopis asjadest, mis on ilmselged, nagu Airflow installimine, teie valitud andmebaas, Celery ja muud teemad, mis on dokumentides kirjeldatud.

Selleks, et me vÔiksime kohe katsetama asuda, olen ma koostanud docker-compose.yml mille sisu on jÀrgmine:

  • KĂ€ivitame enda Airflow: Scheduler, Webserver. Seal töötab ka Flower, et jĂ€lgida Celery ĂŒlesandeid (kuna see on juba pakitud apache/airflow:1.10.10-python3.7, ja me pole selle vastu);
  • PostgreSQL, kuhu Airflow kirjutab oma logifailid (plaanijate andmed, tĂ€itmise statistika jne), ja Celery mĂ€rkis lĂ”petatud ĂŒlesandeid;
  • Redis, mis toimib Celery jaoks ĂŒlesannete vahendajana;
  • Celery worker, mis tĂ”eliselt tegeleb ĂŒlesannete tĂ€itmisega.
  • Katalooge ./dags kasutame meie DAGi kirjeldusfailide salvestamiseks. Need tuvastatakse reaalajas, seega ei ole vaja kogu steki iga vĂ€iksema muudatuse jĂ€rel edastada.

MĂ”nes kohas on kood nĂ€idetes osaliselt esitatud (et mitte teksti ĂŒle koormata), ja mĂ”nes kohas muudetakse seda protsessi kĂ€igus. Terveid ja töötavaid koodinĂ€iteid saab vaadata hoidlas. 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 as a Celery broker
  broker:
    image: redis:6.0.5-alpine

  # DB for the Airflow metadata
  airflow-db:
    image: postgres:10.13-alpine

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

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

  # Main container with Airflow Webserver, Scheduler, Celery Flower
  airflow:
    <<: *airflow-base

    environment:
      <
      -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, will be scaled using `--scale=n`
  worker:
    <<: *airflow-base

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

    depends_on:
      - airflow
      - airflow-db
      - broker

MĂ€rkused:

  • Kompositsiooni koostamisel toetus mulle palju tuntud mudelid. puckel/docker-airflow — kindlasti vaadake. VĂ”ib-olla ei vajagi te elus midagi muud.
  • KĂ”ik Airflow seadistused on saadaval mitte ainult lĂ€bi airflow.cfg, vaid ka keskkonnaparameetrite kaudu (tĂ€nu arendajatele), mille kaudu ma laialdaselt kasutasin.
  • Muidugi ei ole see production-ready: ma ei seadnud konteineritele teadetega ĂŒhendust ja ei aja turvalisusega jamas. Aga miinimum, mis sobib meie katsetamiseks, on loomulikult olemas.
  • Pange tĂ€hele, et:
    • DAG-failide kaust peab olema juurdepÀÀsetav nii ajastajale kui ka töötajatele.
    • Sama kehtib ka kĂ”ikide kolmandate osapoolte raamatukogude kohta — need peavad olema installitud masinatesse, kus on ajastaja ja töötajad.

NĂŒĂŒd on see lihtne:

$ docker-compose up --scale worker=3

PĂ€rast seda, kui kĂ”ik on ĂŒles tĂ”usnud, saab vaadata veebiliideseid:

Peamised mÔisted

Kui te ei saanud aru nendest "DAGidest", siis siin on lĂŒhike sĂ”nastik:

  • Ajastaja — kĂ”ige tĂ€htsam tegelane Airflow's, kes jĂ€lgib, et töötehnikud teevad tööd ja mitte inimene: jĂ€lgib graafikut, vĂ€rskendab DAG-e, kĂ€ivitab ĂŒlesandeid.

    Üldiselt, vanades versioonides olid tal mĂ€luga (ei, mitte amneesia, vaid leketega) probleemid ja konfiguratsioonides jĂ€i isegi alles pĂ€randparameeter. run_duration — tema taaskĂ€ivitamisintervall. Kuid nĂŒĂŒd on kĂ”ik korras.

  • DAG (niisugune "dag") — "suunatud aktsĂŒkliline graaf", kuid selline mÀÀratlemine ĂŒtleb vĂ€he kellelegi, kuigi tegelikult on see konteiner omavahel suhtlevate ĂŒlesannete jaoks (vt allpool) vĂ”i analoog paketile SSIS ja töövoogudele Informatica-s.

    Lisaks DAG-idele vÔivad olla veel subdag-id, kuid me ei jÔua neid tÔenÀoliselt kÀsitleda.

  • DAG Run — initsialiseeritud DAG, millele on antud oma execution_date. Ühe DAG-i DAG-runnid vĂ”ivad kenasti töötada paralleelselt (kui olete oma ĂŒlesanded idempotentseteks teinud).
  • Operator — need on koodijupid, mis vastutavad konkreetse tegevuse tĂ€itmise eest. On kolm tĂŒĂŒpi operaatorit:
    • action, nagu nĂ€iteks meie lemmik PythonOperator, mis on vĂ”imeline tĂ€itma igasugust (kehtivat) Python-koodi;
    • transfer, mis viivad andmeid ĂŒhest kohast teise, nĂ€iteks MsSqlToHiveTransfer;
    • sensor , mis vĂ”imaldab reageerida vĂ”i peatada DAG-i edasise tĂ€itmise mingisuguse sĂŒndmuse toimumiseni. HttpSensor suudab kĂŒsida mÀÀratud lĂ”pp-punkti ja kui ootab vajalikku vastust, alustada ĂŒlekannet GoogleCloudStorageToS3Operator. Uudishimu kĂŒsib: „miks? KĂ”ike saab ju teha otse operaatori sees!” JĂ€rgnevalt, et mitte koormata tĂ¶Ă¶ĂŒlesannete panga riputanud operaatoritega. Sensor kĂ€ivitub, kontrollib ja sureb jĂ€rgmise katse ootamiseks.
  • Ülesanne — kuulutatud operaatorid, sĂ”ltumata tĂŒĂŒbist ja seotud DAGiga, tĂ”stetakse ĂŒlesande tasemele.
  • Ülesande instants — kui peaplanner otsustab, et ĂŒlesanded on valmis heidma vĂ”itlejate töödele (otse kohal, kui me kasutame LocalExecutor vĂ”i kaugnĂ”lva puhul CeleryExecutor), mÀÀrab ta neile konteksti (st muutujate komplekti — tĂ€itmise parameetrid), rakendab kĂ€su vĂ”i pĂ€ringu malle ning paneb need pangale.

Ülesannete genereerimine

Esmalt mÀÀratlege meie DAGi ĂŒldine skeem ja seejĂ€rel sĂŒveneme ĂŒha rohkem detailidesse, sest rakendame mĂ”ningaid ebatavalisi lahendusi.

Nii et kÔige lihtsamas vormis nÀeb selline DAG vÀlja nii:

impordi {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)

Alustame:

  • Esimesena impordime vajalikud teegid ja veel mĂ”ned asjad;
  • sql_server_ds — see on List[namedtuple[str, str]] , mis sisaldab Airflow Connections-i ĂŒhenduste nimesid ja andmebaase, millest me oma tabelit vĂ”tame;
  • dag — meie DAG-i deklareerimine, mis peab kindlasti olema globals(), muidu Airflow ei leia seda. DAG-ile tuleb ka öelda:
    • kuidas seda nimetada orders — see nimi ilmub hiljem veebiliideses,
    • et see hakkab tööle alates 8. juulist keskööst,
    • ja et selle ajakavad on umbes iga 6 tunni jĂ€rel (kĂ”vadele kutidele on siin 'timedelta()' asemel lubatud -string , vĂ€hem kĂ”vadele kutidele - lause nagu cron@daily 0 0 0/6 ? * * *workflow() teeb pĂ”hiosa tööst, aga mitte nĂŒĂŒd. Hetkel lihtsalt vĂ€ljastame meie konteksti logisse.);
  • Ja nĂŒĂŒd lihtne maagia ĂŒlesannete loomise osas: kĂ”nnime allikate jĂ€rgi;
  • algatame
    • , mis tĂ€idab meie tĂŒhi funktsioon
    • initsialiseerime PythonOperator, mis tĂ€idab meie tĂŒhikut Ja nĂŒĂŒd lihtne maagia ĂŒlesannete loomise osas:. Ära unusta anda unikaalne (DAGi piires) ĂŒlesande nimi ning siduda see DAGiga. Lipp provide_context toob omakorda funktsiooni tĂ€iendavad argumendid, mille me ettevaatlikult kogume **context.

Siinkohal lÔpetame. Mida me oleme saanud:

  • uus DAG veebiliideses,
  • kahesaja viiskĂŒmmend ĂŒlesannet, mis töötavad paralleelselt (kui Airflow, Celery ja serverite seaded seda lubavad).

Noh, peaaegu saime.

Apache Airflow: teeme ETL lihtsamaks
Kes paigaldab sÔltuvused?

Kogu seda asja lihtsustada, lĂŒkkasin ma docker-compose.yml töötlemisse requirements.txt kĂ”ikidesse nodidesse.

NĂŒĂŒd lĂ€heb asi kiireks:

Apache Airflow: teeme ETL lihtsamaks

Hallid ruudud — ĂŒlesande instantsid, mida planeerija on töötlenud.

Ootame veidi, ĂŒlesanded haaravad töötajad:

Apache Airflow: teeme ETL lihtsamaks

Rohelised, loomulikult, — edukalt lĂ€bitud. Punased — mitte nii edukalt.

Muide, meie tootmises ei ole mingit kausta ./dags, mis sĂŒnkroniseerub masinate vahel — kĂ”ik DAGid asuvad git meie Gitlabis, ja Gitlab CI seab uuendused masinatesse, kui need on kokku pandud master.

Natuke Flowerist

Kuni töötajad töötlevad meie tĂŒhje ĂŒlesandeid, tuletame meelde veel ĂŒhe tööriista, mis vĂ”ib meile mĂ”ningaid asju nĂ€idata — Flower.

Esimene leht, kus on kokkuvÔtete info töötajate nodide kohta:

Apache Airflow: teeme ETL lihtsamaks

KĂ”ige sisukam leht ĂŒlesannetega, mis on töös:

Apache Airflow: teeme ETL lihtsamaks

KÔige igavam leht meie maaklerite seisuga:

Apache Airflow: teeme ETL lihtsamaks

KĂ”ige silmatorkavam leht — ĂŒlesannete seisugraafik ja nende tĂ€itmise ajad:

Apache Airflow: teeme ETL lihtsamaks

Laadime alla, mis ei ole laaditud

Nii, kĂ”ik ĂŒlesanded on töös, saame kannatanud vĂ€lja viia.

Apache Airflow: teeme ETL lihtsamaks

Ja kannatanuid oli pĂ€ris mitu — erinevatel pĂ”hjustel. Kui Airflowd Ă”igesti kasutada, siis need ruudud nĂ€itavad, et andmed pole kindlasti kohale jĂ”udnud.

Tuleb vaadata logi ja taaskĂ€ivitada kukkunud ĂŒlesandeinstantsid.

Klikkides mÔnel ruudul, nÀeme meie jaoks saadaval olevaid tegevusi:

Apache Airflow: teeme ETL lihtsamaks

Saame vĂ”tta ja teha kukkunutele Clear. See tĂ€hendab, et unustame, et seal midagi takerdus, ja sama ĂŒlesandeinstants suundub planeerijale.

Apache Airflow: teeme ETL lihtsamaks

Selge, et teiste punaste ruutudega hiirega nii kĂ€ituda ei ole kuigi inetu — mitte seda me Airflowlt ootame. Loomulikult on meil massihĂ€vitusrelv: Browse/Task Instances

Apache Airflow: teeme ETL lihtsamaks

Valime kÔik korraga ja nullime, vajutame Ôiget punkti:

Apache Airflow: teeme ETL lihtsamaks

PÀrast puhastamist nÀevad meie taksod vÀlja nii (nad juba ootavad, millal planeerija nad plaanib):

Apache Airflow: teeme ETL lihtsamaks

Ühendused, hooikud ja muud muutujad

On viimane aeg vaadata jÀrgmist DAG-i, 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("""Tere, head inimesed, raportid on uuendatud"""),
    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, Ă€rka ĂŒles, me kukutasime {{ 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]

Kas oleme kĂ”ik kunagi aru saanud aruandeid uuendama? Siin see on jĂ€lle: meil on loetelu allikatest, kust andmed vĂ”tta; on loetelu, kuhu need panna; Ă€rme unustame signaale anda, kui kĂ”ik on sĂŒndinud vĂ”i purunenud (aga see ei kehti meie kohta, eks?).

Vaadakem uuesti faili ja vaatame uusi mÔistetavaid asju:

  • from commons.operators import TelegramBotSendMessage — mis takistab meid oma operaatorite loomisel, ja me kasutasime seda Ă€ra, tehes vĂ€ikese mĂ€hise sĂ”numite saatmiseks Razblokirovanny'sse. (Selle operaatori kohta rÀÀgime veel allpool);
  • default_args={} — DAG vĂ”ib jagada samu argumente kĂ”igile oma operaatoritele;
  • to='{{ var.value.all_the_kings_men }}' — vĂ€li to meil ei ole seda kĂ”vakooditud, vaid dĂŒnaamiliselt genereeritav Jinja ja muutujaga, kus on e-posti aadresside loend, mille ma ettevaatlikult panin Admin/Variables;
  • trigger_rule=TriggerRule.ALL_SUCCESS — operaatori kĂ€ivitamise tingimus. Meie juhul saadetakse kiri juhtidele ainult juhul, kui kĂ”ik sĂ”ltuvused on töötanud edu;
  • tg_bot_conn_id='tg_main' — argumendid conn_id vĂ”tavad endasse ĂŒhenduste identifikaatorid, mille me loome Admin/Connections;
  • trigger_rule=TriggerRule.ONE_FAILED — sĂ”numid Telegramis saadetakse ainult siis, kui on ebaĂ”nnestunud ĂŒlesandeid;
  • task_concurrency=1 — keelame ĂŒhe ĂŒlesande mitme taski instantsi samaaegse kĂ€itamise. Vastasel juhul saame mitu ĂŒlesande instantsi, VerticaOperator (mis vaatavad ĂŒhte tabelit);
  • report_update >> [email, tg] — kĂ”ik VerticaOperator koondatakse e-kirjade ja sĂ”numite saatmiseks, nagu jĂ€rgneb:
    Apache Airflow: teeme ETL lihtsamaks

    Kuna aga teavitajate operaatoritel on erinevad kĂ€ivitamistingimused, töötab vaid ĂŒks. Tree View'is nĂ€eb see vĂ€lja vĂ€hem lĂ€bipaistev:
    Apache Airflow: teeme ETL lihtsamaks

RÀÀgin natuke makrode ja nende sĂ”prade — muutujate.

Makrod on Jinja-sisu kohandajad, mis vÔivad asetada erinevat kasulikku teavet operaatorite argumentidesse. NÀiteks nii:

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

{{ ds }} laiendab konteksti muutujate sisu execution_date vormingus YYYY-MM-DD: 2020-07-14. Parim on see, et konteksti muutujad on kinnitatud kindla taski instantsi (ruudu) kĂŒlge Tree View'is, ja kui instantsi uuesti kĂ€ivitada, lahetakse kohandajad samadele vÀÀrtustele.

MÀÀratud vÀÀrtusi saab vaadata igas taski instantsis nupul Rendered. Nii on e-kirjade saatmise ĂŒlesande puhul:

Apache Airflow: teeme ETL lihtsamaks

Ja nii on sĂ”numi saatmise ĂŒlesande puhul:

Apache Airflow: teeme ETL lihtsamaks

Viimane koodiversioon sisaldab kÔiki integreeritud makrode loendit, mis on saadaval siin: Makrode viide

Lisaks saame pluginade abil luua oma makrosid, kuid sellest rÀÀgime hiljem.

Eeldefineeritud asjade kÔrval saame sisestada ka oma muutuja vÀÀrtusi (nagu olen juba koodis teinud). Loome mÔned: Admin/Variables paar asja:

Apache Airflow: teeme ETL lihtsamaks

Nii, nĂŒĂŒd saab kasutada:

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

VÀÀrtuses vÔib olla skalaarsed andmed vÔi isegi JSON. JSON-i puhul:

bot_config

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

kasutame lihtsalt Ôige vÔtme teed: {{ var.json.bot_config.bot.token }}.

Ütlen lihtsalt ĂŒhe sĂ”na ja nĂ€itan ĂŒhte ekraanipilti seoses ĂŒhendustega. Siin on kĂ”ik lihtne: lehe peal Admin/Connections loome ĂŒhenduse, salvestame sinna oma kasutajanimesid/paroolid ja spetsiifilisemad parameetrid. Nii:

Apache Airflow: teeme ETL lihtsamaks

Paroolid saab krĂŒpteerida (tĂ€iendavalt, kui vĂ”rrelda vaikeseadetega) vĂ”i ei pea ka ĂŒhenduse tĂŒĂŒpi mainima (nagu ma tegin) tg_main) — asi on selles, et Airflow mudelites on tĂŒĂŒpide nimekiri sisse ehitatud ja seda ei saa muudetud allikateta laiendada (kui ma midagi valesti otsisin — palun parandage mind), kuid me ei takista midagi, et saada akrediteeringud lihtsalt nime jĂ€rgi.

Samuti saab luua mitu ĂŒhendust sama nimega: sellisel juhul meetod BaseHook.get_connection(), mis toob meile ĂŒhendused nime jĂ€rgi, annab juhusliku mitme sarnase seas (loogilisem oleks olnud kasutada Round Robin'i, kuid jĂ€tame selle Airflow arendajate sĂŒdametunnistusele).

Muutujad ja ĂŒhendused on kindlasti suurepĂ€rased tööriistad, kuid on oluline hoida tasakaalu: millised osad teie voogudest te hoiate koodis ja milliseid — annate Airflow’le sĂ€ilitamiseks. Ühelt poolt vĂ”ib vÀÀrtuse, nĂ€iteks postkasti, kiirelt muutmine UI kaudu olla mugav. Teiselt poolt on see siiski tagasipöördumine hiireklikkide juurde, millest me (mina) tahtsime vabaneda.

Ühendustega töötamine on ĂŒks hook'ideĂŒlesandeid. Üldiselt on Airflow hook'id side punktid kolmandate teenuste ja raamatukogudega. NĂ€iteks, JiraHook avatakse meie jaoks klient, et suhelda Jira'ga (saame ĂŒlesandeid edasi-tagasi liigutada), ja SambaHook abil saab pushida kohaliku faili smb-punkti.

Kohandatud operaatori analĂŒĂŒs

Ja me oleme lÀhenenud sellele, et vaadata, kuidas on tehtud TelegramBotSendMessage

Kood commons/operators.py isegi operaatori juurde:

from typing import Union

from airflow.operators import BaseOperator

from commons.hooks import TelegramBotHook, TelegramBot

class TelegramBotSendMessage(BaseOperator):
    """Saada sÔnum chat_id-le, kasutades TelegramBotHook'i

    NĂ€ide:
        >>> 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 }} ebaÔnnestus :(',
        ...     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'Saada "{self.message}" chat' + f'{self.chat_id}')
        self.client.send_message(chat_id=self.chat_id,
                                 message=self.message)

Siin, nagu kÔik muu Airflow's, on kÔik vÀga lihtne:

  • Oleme pĂ€rinud BaseOperator, mis rakendab ĂŒsna palju Airflow'le spetsiifilisi asju (vaadake rahus)
  • Oleme kuulutanud vĂ€lja vĂ€ljad template_fields, kus Jinja otsib makrosid töötlemiseks.
  • Oleme korraldanud Ă”iged argumendid __init__(), mÀÀrasid vaikimisi, kus vaja.
  • EeljĂ€tku initsialiseerimist ei unustatud.
  • Avatud vastav hook. TelegramBotHook, saime temalt klient-objekti.
  • Ülekirjutasime meetodi. BaseOperator.execute(), mida Airflow kutsub, kui on aega operaatori kĂ€ivitamiseks — selles realiseerime pĂ”hitegevuse, unustamata logida. (Logime, muide, kohe. stdout ja stderr — Airflow kĂ”ik pĂŒĂŒab kinni, pakib ilusti kokku ja paigutab Ă”igesse kohta.)

Vaadakem, mis meil on failis commons/hooks.py. Faili esimene osa, mis sisaldab hooki:

from typing import Union

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

class TelegramBotHook(BaseHook):
    """Telegram Bot API hook

    MĂ€rkus: lisage ĂŒhendus tĂŒhja Conn Type'iga ja Ă€rge unustage
    tÀita 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

Ma isegi ei tea, mida siin seletada, lihtsalt toon vÀlja mÔned olulised punktid:

  • Kasutame pĂ€randit, mĂ”tleme argumentide ĂŒle — enamikul juhtudel on see ĂŒks: conn_id;
  • Üksikute standardmeetodite ĂŒletamine: piirdun get_conn(), kus ma saan ĂŒhenduse parameetrid nime jĂ€rgi ja lihtsalt tĂ”mban sektsiooni extra (see vĂ€li on JSON-i jaoks), kuhu ma (kui ma ise juba selgitasin!) panin Telegrami boti tokeni: {"bot_token": "YOuRAwEsomeBOtToKen"}.
  • Loome meie TelegramBot, edastades talle konkreetse tokeni.

Ja ongi kĂ”ik. Saame kliendi ĂŒlesandest lĂ€bi TelegramBotHook().clent vĂ”i TelegramBotHook().get_conn().

Ja teises osas faile, kus ma teen mikropakkumise Telegram REST API jaoks, et mitte vedada sama python-telegram-bot ĂŒhe meetodi jaoks sendMessage.

class TelegramBot:
    """Telegram Bot API wrapper

    Examples:
        >>> TelegramBot('YOuRAwEsomeBOtToKen', '@myprettydebugchat').send_message('Hi, darling')
        >>> TelegramBot('YOuRAwEsomeBOtToKen').send_message('Hi, darling', 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))

Õige tee on kĂ”ik see kokku panna: TelegramBotSendMessage, TelegramBotHook, TelegramBot — plugina, panna avalikku hoidlasse ja jagada avatud lĂ€htekoodiga.

Kuna me kÔike seda uurisime, olid meie raporti vÀrskendused edukalt ilmu t seekord kokku kukkunud ja saatsid mulle kanalis vea teate. Pean minema kontrollima, mis seekord valesti lÀks


Apache Airflow: teeme ETL lihtsamaks
Meie dags on katki! Ja ei olnud just seda, mida me ootasime? Just nimelt!

Kas hakkad valama?

Kas tunnete, et jĂ€tan midagi tĂ€helepanuta? Tundub, et lubasin andmed SQL Serverist Vertica'le ĂŒle kanda, ja nĂŒĂŒd lĂ€ksin teemadest kĂ”rvale, kurat!

Kuritegu oli tahtlik, ma pidin teile selgitama mĂ”ningaid termineid. NĂŒĂŒd saame liikuda edasi.

Meie plaan oli jÀrgmine:

  1. Luua dag
  2. Genererida ĂŒlesanded
  3. Vaadata, kui ilus kÔik vÀlja nÀeb
  4. Kinnituselamistele jagada sessiooninumbreid
  5. Andmete vÔtmine SQL Serverist
  6. Andmete paigutamine Verticasse
  7. Statistika kogumine

Nii et selle kÔik kÀivitamiseks tegin vÀikese tÀienduse meie 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

Seal me tÔstame:

  • Vertica kui host dwh kĂ”ige vaikimisi seadistustega,
  • kolm SQL Serveri instantsi,
  • tĂ€idame andmebaasid viimaste andmetega (Ă€rge lĂŒkake vaadata mssql_init.py!)

KÀivitame kogu selle hea veidi keerulisema kÀsu abil kui eelmisel korral:

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

Mida meie imetore juhuslik generaator genereeris, saab kasutada punkti kaudu Andmete profiilimine/Ad Hoc pÀring:

Apache Airflow: teeme ETL lihtsamaks
Peamine, et seda analĂŒĂŒtikutele ei nĂ€idata

Detailidesse laskuda ETL-seansid ei hakka, seal on kĂ”ik triviaalne: loome andmebaasi, sinna tabeli, mĂ€hkime kĂ”ik konteksti haldajasse ja nĂŒĂŒd teeme nii:

with Session(task_name) as session:
    print('Load', session.id, 'started')

    # Laadimise töövoog
    ...

    session.successful = True
    session.loaded_rows = 15

session.py

from sys import stderr

class Session:
    """ETL töövoo sessioon

    NĂ€ide:
        with Session(task_name) as session:
            print(session.id)
            session.successful = True
            session.loaded_rows = 15
            session.comment = 'Tore töö'
    """

    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 = """
            LOONI TABEL, KUI EI OLE JUBA olemas sessioonid (
                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, 'avatud')
        return self

    def close(self):
        if not self._id:
            raise SessionClosedError('Sessioon ei ole avatud')
        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, 'suletud',
              ', edukas: ', self.successful,
              ', Laetud: ', self.loaded_rows,
              ', kommentaar:', self.comment)

class SessionError(Exception):
    pass

class SessionClosedError(SessionError):
    pass

On aeg teha meie andmete vÔtmisest meie poolest sada tabelit. Teeme seda lihtsate ridade abil:

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. Kasutame Airflow'd, et saada pymssql-ĂŒhendust
  2. Lisame pĂ€ringusse kuupĂ€eva piirangu — funktsioon lisab selle ĆĄabloonide kaudu.
  3. Saadame meie pĂ€ringu pandas, mis toob meile DataFrame — see tuleb meile kasuks edaspidi.

Kasutame asendust {dt} pÀringu parameetri asemel %s mitte sellepÀrast, et ma oleks kuri Buratino, vaid sellepÀrast, et pandas ei oska hakkama saada pymssql ja annab viimasele edasi params: List, kuigi see vÀga tahab tuple.
Olge tĂ€helepanelik, et arendaja pymssql otsustas seda enam mitte toetada, ja on viimane aeg minna ĂŒle pyodbc.

Vaadakem, millised argumendid Airflow meie funktsioonidesse sisestas:

Apache Airflow: teeme ETL lihtsamaks

Kui andmeid polnud, siis pole mĂ”tet edasi minna. Kuid arvestada, et ĂŒleslaadimine oli edukas, on ka kummaline. Kuid see ei ole ka viga. Ah, mida teha?! Niimoodi:

if df.empty:
    raise AirflowSkipException('No rows to load')

AirflowSkipException Airflow ĂŒtleb, et viga ei ole, ja ĂŒlesanne jĂ€etakse vahele. Liideses ei ole rohelise ega punase ruudu asemel roosa vĂ€rv.

Lisame oma andmetele mitu veergu:

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

Konkreetsemalt:

  • Andmebaas, kust me tellimused saime,
  • Meie laadimise seansi identifikaator (see erineb kĂ”igil ĂŒlesannetel),
  • Ressursi ja tellimuse identifikaatorist saadud hash — et lĂ”pp-andmebaasis (kus kĂ”ik kogutakse ĂŒhte tabelisse) oleks ainulaadne tellimuse identifikaator.

JÀÀnud on eelviimane samm: laadida kĂ”ik Verticasse. Ja, kummalisel kombel, ĂŒks kĂ”ige efektsematest viisidest selleks on CSV kaudu!

# 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. Loome spetsiaalse vastuvÔtja StringIO.
  2. pandas mis kenasti kogub meie DataFrame kujul CSV-ridu.
  3. Avame ĂŒhenduse meie lemmik Vertica hookiga.
  4. Ja nĂŒĂŒd kasutame copy() saadame meie andmed otse Verticasse!

Draiverist saame, kui palju ridu laaditi, ja ĂŒtleme seansi juhile, et kĂ”ik on OK:

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

Ja kÔik.

Tootmises loome sihttabeli kÀsitsi. Siin julgen lubada vÀikest automatiseerimist:

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)

Ma olen abiks VerticaOperator() loome andmebaasi skeemi ja tabeli (kui neid pole, muidugi). Peamine on Ôigesti seada sÔltuvused:

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

Teeme kokkuvÔtteid

— Noh, — ĂŒtles hiir, — ei ole tĂ”si, et nĂŒĂŒd
Sa veendusid, et metsas olen mina kÔige hirmsam metsaline?

Julia Donaldson, „Gruffalo“

Ma arvan, et kui me oma kolleegidega korraldaksime konkurentsi: kes suudab kiiremini koostada ja kĂ€ivitada ETL-protsessi nullist: nemad oma SSIS ja hiirega ning mina Airflow’ga
 Ja seejĂ€rel vĂ”rdleksime hooldamise mugavust
 Oh, arvan, et nĂ”ustute, et ma tĂ”enĂ€oliselt ĂŒletan neid igal rindel!

Kui vĂ”tta asja veidi tĂ”sisemalt, siis Apache Airflow — tĂ€nu protsesside kirjeldamisele programmeerimiskoodina — tegi minu töö kaugelt mugavamaks ja meeldivamaks.

Tema piiramatud laienemisvĂ”imalused: nii pistikprogrammide osas kui ka skaleeritavuse poolest — annavad teile vĂ”imaluse rakendada Airflow’d praktiliselt igas valdkonnas: olgu see siis andmete kogumise, töötlemise ja ettevalmistamise tĂ€islĂŒka vĂ”i rakettide kĂ€ivitamine (Marsile, muidugi).

LÔpposa, viidatud ja informatiivne

KĂŒhvlid, mis me teie eest kogusime

  • start_date. Jah, see on juba kohalike meemide teemade ring. Peamise dĂŒnaamika argumendi kaudu start_date lĂ€bivad kĂ”ik. LĂŒhidalt, kui mÀÀrata start_date praegune kuupĂ€ev ja schedule_interval — ĂŒks pĂ€ev, siis DAG kĂ€ivitub homme mitte varem.
    start_date = datetime(2020, 7, 7, 0, 1, 2)

    Ja rohkem ei mingeid probleeme.

    Sellega on seotud ka ĂŒks teine tĂ€itmisviga: Task is missing the start_date parameter, mis enamasti tĂ€hendab, et unustasite seondada DAG operaatoriga.

  • KĂ”ik ĂŒhel masinal. Jah, nii Airflow andmebaasid kui ka meie kattekiht, veebiserver, ajakavaja ja töötajad. Ja see töötas isegi. Kuid aja jooksul kasvas teenuste ĂŒlesannete arv ning kui PostgreSQL hakkas andma vastuse indeksi jĂ€rgi 20 ms asemel 5 ms, siis me vĂ”tsime selle ja viibisime Ă€ra.
  • LocalExecutor. Jah, me kasutame seda endiselt ja oleme jĂ”udnud ÀÀre poole. LocalExecutor on seni olnud piisav, kuid nĂŒĂŒd on aeg vĂ€hemalt ĂŒhe töötaja vĂ”rra laieneda ja tuleb pingutada, et ĂŒle minna CeleryExecutorile. Kuna sellega saab töötada ka ĂŒhel masinal, ei ole miski takistuseks kasutada Celery't isegi serveris, mis "loomulikult ei lĂ€he kunagi tootmisse, ausĂ”na!"
  • Kasutamata sisseehitatud tööriistu:
    • Ühendused teenuste autentimisandmete salvestamiseks,
    • SLA puudumised ĂŒlesannete jaoks, mis ei töötanud Ă”igel ajal,
    • XCom metaandmete vahetamiseks (ma ĂŒtlesin metaandmed!) ĂŒlesannete vahel DAG'is.
  • E-kirjade liialdamine. Mis siin muud öelda? Oleme seadnud hoiatused kĂ”igi ebaĂ”nnestunud ĂŒlesannete korduste kohta. NĂŒĂŒd on mu töö Gmailis >90k kirja Airflow'ilt ja e-kirjade veebirakendus keeldub vĂ”tma ja kustutama rohkem kui 100 korraga.

Rohkem takistusi: Apache Airflow'i lÔksud

Veel suurema automatiseerimise vahendid

Kuna soovime veelgi rohkem mÔtlema hakata, mitte kÀtega töötada, on Airflow meie jaoks valmistanud jÀrgmise:

  • REST API — ta jĂ€tkuvalt omab Ekspressi staatust, mis ei takista selle tööleminekut. Selle abil on vĂ”imalus mitte ainult saada teavet DAGide ja ĂŒlesannete kohta, vaid ka peatada/algatada DAGi, luua DAGi jooksu vĂ”i hulki.
  • CLI — kĂ€surealt on kergesti kĂ€ttesaadavad paljud tööriistad, mis ei ole mitte ainult ebamugavad WebUI's, vaid ei ole seal ĂŒldse saadaval. NĂ€iteks:
    • backfill on vajalik ĂŒlesannete instantside uuesti kĂ€ivitamiseks.
      NĂ€iteks kui analĂŒĂŒtikud tulevad ja ĂŒtlevad: «Teie andmeanalĂŒĂŒs on halb ja vajab parandamist 1. kuni 13. jaanuarini! Parandage, parandage, parandage, parandage!» Siis vĂ”tad ja teed:
      airflow backfill -s '2020-01-01' -e '2020-01-13' orders
    • Andmebaasi hooldus: initdb, resetdb, upgradedb, checkdb.
    • run, mis vĂ”imaldab kĂ€ivitada ĂŒhe ĂŒlesande instantsi, jĂ€ttes kĂ”ik sĂ”ltuvused tĂ€helepanuta. Veelgi enam, seda on vĂ”imalik kĂ€ivitada lĂ€bi LocalExecutor, isegi kui teil on Celery-klaster.
    • Umbes sama asja teeb test, ainult et ei kirjuta midagi andmebaasi.
    • connections vĂ”imaldab massiliselt luua ĂŒhendusi shell'ist.
  • Python API — ĂŒsna hardcore viis suhtlemiseks, mis on mĂ”eldud pistikprogrammide jaoks, mitte kĂ€sitsi tööks. Aga kes meid kĂŒll takistab minemast /home/airflow/dags, kĂ€ivitada ipython ja hakkama saama? NĂ€iteks on vĂ”imalik kĂ”ik ĂŒhendused eksportida sellise koodiga:
    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)
  • Airflow'i metaandmebaasiga ĂŒhendamine. Soovitan sinna kirjutada mitte, kuid erinevate spetsiifiliste mÔÔdikute jaoks ĂŒlesannete olekute kĂ€tte saamine on kindlasti kiirem ja lihtsam kui lĂ€bi ĂŒkskĂ”ik millise API.

    Ütleme, et kĂ”ik meie ĂŒlesanded ei ole idempotentsed, ja need vĂ”ivad mĂ”nikord ebaĂ”nnestuda, ja see on normaalne. Kuid mitu kokkuvarisemist on juba kahtlane ning tuleks kontrollida.

    Ole ettevaatlik, SQL!

    WITH last_executions AS (
    SELECT
        task_id,
        dag_id,
        execution_date,
        state,
            row_number()
            OVER (
                PARTITION BY task_id, dag_id
                ORDER BY execution_date DESC) AS rn
    FROM public.task_instance
    WHERE
        execution_date > now() - INTERVAL '2' DAY
    ),
    failed AS (
        SELECT
            task_id,
            dag_id,
            execution_date,
            state,
            CASE WHEN rn = row_number() OVER (
                PARTITION BY task_id, dag_id
                ORDER BY execution_date DESC)
                     THEN TRUE END AS last_fail_seq
        FROM last_executions
        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 failed
    GROUP BY
        task_id,
        dag_id
    HAVING
        count(last_fail_seq) > 0

Lingid

Ja ja, ja esiteks kĂŒmme linki Google'i otsingust, mis viivad minu Airflow kaustade sisule.

Ja lingid, mis on artiklis kasutusel:

Allikas: habr.com

Osta usaldusvÀÀrne veebihosting DDoS kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne veebihosting DDoS kaitsega, VPS VDS serverid | ProHoster