Apache Airflow: lihtsustame ETL-i

Tere, mina olen Dmitri Logvinenko — andmete insener ettevĂ”tte «Vezёт» analĂŒĂŒsiosakonnas.

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

Ja jah, ma mitte ainult ei rÀÀgi, vaid ka nÀitan: programmis on palju koodi, ekraanipilte ja soovitusi.

Apache Airflow: lihtsustame ETL-i
Mida tavaliselt nÀed, kui googled sÔna Airflow / Wikimedia Commons

Sisukord

Sissejuhatus

Apache Airflow — ta on nagu Django:

  • kirjutatud Pythonis,
  • on suurepĂ€rane adminpaneel,
  • piiramatu laiendatavus,

— ainult parem, ja loodud hoopis teiste eesmĂ€rkide saavutamiseks, nimelt (kui on kirjutatud enne katti):

  • ĂŒlesannete kĂ€ivitamine ja jĂ€lgimine piiramatul arvul masinatel (nii palju kui teie lubab Celery/Kubernetes ja teie sĂŒdametunnistus)
  • dĂŒnaamilise töövoo genereerimine vĂ€ga lihtsalt kirjutatava ja tajutava Python-koodi kaudu
  • ja vĂ”imalus siduda omavahel ĂŒkskĂ”ik millised andmebaasid ja API-d nii valmis komponentide kui ka isetehtud pluginatega (mida on ÀÀrmiselt lihtne teha).

Me kasutame Apache Airflow'i nii:

  • kogume andmeid erinevatest allikatest (palju SQL Serveri ja PostgreSQL instantsse, erinevaid API-sid rakenduste mÔÔdikute jaoks, isegi 1C) DWH-sse ja ODS-i (meie puhul on need Vertica ja Clickhouse).
  • nagu arenenud cron, mis kĂ€ivitab andmete konsolideerimise protsesse ODS-is ja jĂ€lgib nende hooldust.

Kuni hiljutise ajani katab meie vajadusi ĂŒks vĂ€ike server 32 tuuma ja 50 GB RAM-iga. Airflow'is töötab sel juhul:

  • ĂŒle 200 DAG-i (nende töövood, kuhu oleme ĂŒlesandeid tĂ€itnud),
  • igaĂŒhes keskmiselt 70 ĂŒlesannet,
  • selle headuse kĂ€ivitamine (ka keskmiselt) korra tunni jooksul.

Ja sellest, kuidas me laienesime, kirjutan ma allpool, kuid praegu mÀÀratlegem ĂŒber-ĂŒlesanne, mida me lahendama hakkame:

On each of the three source SQL Servers, there are 50 databases — instances of one project, so their structure is identical (almost everywhere, muah-ha-ha), which means each has a table called Orders (after all, a table with that name can fit into any business). We pull data, adding metadata fields (source server, source database, ETL task identifier) and naively drop them into, say, Vertica.

LĂ€hme!

Peamine, praktiline (ja natuke teoreetiline) osa

Miks see meile (ja teile) vajalik on

When the trees were tall, and I was just a simple SQL-worker in a Russian retail company, we rolled out ETL processes aka data streams using the two tools available to us:

  • Informatica Power Center — an extremely sprawling system, highly productive, with its own hardware, its own versioning. I used maybe 1% of its capabilities. Why? Well, first, this interface is from the noughties and mentally weighed down on us. Secondly, this thing is designed for incredibly complex processes, fierce component reuse, and other very-important-enterprise-features. Let's not even mention that it costs as much as the wing of an Airbus A380/year.

    Caution, the screenshot may hurt people under 30 a bit.

    Apache Airflow: lihtsustame ETL-i

  • SQL Server Integration Server — we used this fellow in our internal project streams. Well, in fact: we already use SQL Server, and not using its ETL tools would be unreasonable. Everything about it is good: the interface is beautiful, and the execution reports
 But that’s not why we love software products, oh no. We can version it dtsx (which is an XML with nodes mixed when saving); but what's the use? And to create a task package that will move a hundred tables from one server to another? You’d drop your index finger on the mouse button before moving twenty. But it definitely looks more stylish:

    Apache Airflow: lihtsustame ETL-i

We were definitely looking for ways out. It even almost came to a self-made SSIS package generator



 and then I found a new job. And on that job, I encountered Apache Airflow.

When I learned that describing ETL processes is just simple Python code, I almost danced with joy. This is how data streams underwent versioning and diffusion, and pouring tables with a uniform structure from a hundred databases into one target became a task for Python code on one and a half to two 13" screens.

Kogume klastrit

Ärge korraldage lastaeda ja rÀÀkige tĂ€iesti ilmsetest asjadest, nagu Airflow'i, teie valitud andmebaasi, Celery ja muude asjade seadistamisest, millest dokumentatsioonis rÀÀgitakse.

Kuna me saame kohe katsetama hakata, olen ma koostanud docker-compose.yml kuna:

  • KĂ€ivitame tegelikult Airflow: Scheduler, Webserver. Seal kĂ€ivitub ka Flower Celery ĂŒlesannete jĂ€lgimiseks (sest see on juba lĂŒkatud apache/airflow:1.10.10-python3.7, ja me ei ole selle vastu);
  • PostgreSQL, kuhu Airflow kirjutab oma haldusandmed (planeerija andmed, tĂ€itmise statistika jne), ja Celery tĂ€histab lĂ”petatud ĂŒlesandeid;
  • Redis, mis toimib Celery ĂŒlesannete vahendajana;
  • Celery worker, mis hakkab ĂŒlesandeid tegelikult tĂ€itma.
  • Kausta . /dags salvestame meie DAGide kirjeldusfailid. Need tuuakse automaatselt ĂŒles, seega pole vaja iga kord kogu steki tĂ€iesti ĂŒles tĂ”sta.

Kohati on kood nÀidetes osaliselt toodud (et mitte teksti ummistada) ja mÔned kohad muudetakse protsessi kÀigus. TÀielikke töötavaid koodinÀiteid saab vaadata hoidlas. https://github.com/dm-logv/airflow-tutorial.

docker-compose.yml

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

    environment:
      <<: *airflow-config

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

    depends_on:
      - airflow
      - airflow-db
      - broker

MĂ€rkused:

  • Ma kompositsiooni koostamisel toitusin peamiselt tuntud pildist puckel\/docker-airflow – kindlasti vaadake. VĂ”ib-olla ei vaja te elus midagi muud.
  • KĂ”ik Airflow seaded on saadaval mitte ainult lĂ€bi airflow.cfg, vaid ka keskkonnamuutujate kaudu (tĂ€nu arendajatele), mida ma hĂ€bivÀÀrselt kasutasin.
  • Muidugi ei ole see production ready: ma ei seadnud konteineritele heartbeate, ma ei vaevanud end turvalisusega. Kuid miinimum, mis sobib meie katsetele, on loodud.
  • Pange tĂ€hele, et:
    • DAGide kaust peab olema kergesti ligipÀÀsetav nii ajakava koostajale kui ka töötajatele.
    • Sama kehtib kĂ”igi kolmandate osapoolte teekide kohta — need peavad kĂ”ik olema installitud ajakava koostaja ja töötajate masinatesse.

NĂŒĂŒd aga lihtsalt:

$ docker-compose up --scale worker=3

Kui kĂ”ik on kĂ€ima lĂŒkatud, saab veebiliideseid vaatama minna:

PÔhimÔisted

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

  • Scheduler — peamine tegelane Airflow's, kes kontrollib, et robotid töötaksid, mitte inimene: jĂ€lgib ajakava, uuendab dag'e, kĂ€ivitab ĂŒlesandeid.

    Tegelikult olid vanade versioonide puhul tal mĂ€luprobleemid (ei, mitte amneesia, vaid lekked) ja konfiguratsioonides jĂ€i isegi alles ĂŒks legassi parameeter run_duration — tema taaskĂ€ivitamise intervall. Aga nĂŒĂŒd on kĂ”ik korras.

  • DAG (ehk "dag") — "suunatud ahelate graaf", kuid see mÀÀratlemine ei ĂŒtle kellelegi suurt midagi, aga sisuliselt on see konteiner omavahel suhtlevatele ĂŒlesannetele (vt allpool) vĂ”i analoog Paketile SSIS-is ja Töökargile Informatica-s.

    Lisaks dag'idele vÔivad olla ka subdag'id, aga tÔenÀoliselt me needeni ei jÔua.

  • DAG Run — initsialiseeritud dag, millele on antud oma execution_date. Ühe dag'i dag'runnid vĂ”ivad tĂ€iesti töötada paralleelselt (kui olete muidugi teinud oma ĂŒlesanded idempotentseteks).
  • Operator — see on koodi tĂŒkk, mis vastutab mingi konkreetse toimingu tĂ€itmise eest. On kolm tĂŒĂŒpi operaatorit:
    • tegevus, nagu nĂ€iteks meie lemmik PythonOperator, mis on vĂ”imeline tĂ€itma mis tahes (kehtivat) Python-koodi;
    • ĂŒlekandmine, mis viivad andmeid ĂŒhelt poolt teisele, nĂ€iteks MsSqlToHiveTransfer;
    • sensor lubab reageerida vĂ”i peatada edasise dag'i tĂ€itmise, kuni mingi sĂŒndmus ei ole toimunud. HttpSensor vĂ”ib kĂŒsida mÀÀratud lĂ”pp-punkti, ja kui ootab vajaliku vastuse, kĂ€ivitab körvĂŒlekande GoogleCloudStorageToS3Operator. Uudishimu kĂŒsib: "miks? Kas poleks lihtsam teha kordusi otse operaatoris!" Ja siis, et mitte koormata ĂŒlesannete kogumit seisatavate operaatoritega. Sensor kĂ€ivitub, kontrollib ja sureb kuni jĂ€rgmise katseni.
  • Task — vĂ€ljakuulutatud operaatorid, olenemata tĂŒĂŒbist, ja dag'iga seotud ĂŒlesannete mÀÀramine tĂ”stetakse ĂŒlesande tasemele.
  • Task instance — kui peaplaneerija otsustas, et ĂŒlesanded on valmis lĂ€hetamiseks tote jooksutavatele töötajatele (otse kohapeal, kui kasutame LocalExecutor vĂ”i kaug-noodi puhul CeleryExecutor), mÀÀrab ta neile konteksti (st muutuja-sĂ€tte komplekti), rakendab kĂ€skude vĂ”i pĂ€ringute malle ja kogub need kokku.

Genereerime ĂŒlesandeid

Esmalt mÀÀratleme meie dag'i ĂŒldise skeemi, seejĂ€rel sĂŒveneme jĂ€rjest enam detailidesse, sest rakendame mitmeid mitte triviaalsetes lahendustes.

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

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)

Hakkame siis ĂŒheskoos vĂ€lja selgitama:

  • Esmalt impordime vajalikud teegid ja veel mĂ”ned muud asjad;
  • sql_server_ds — see on List[namedtuple[str, str]] Airflow Connections'i ĂŒhenduste nimed ja andmebaasid, millest me meie tabelit toome;
  • dag on meie DAG'i deklaratsioon, mis peab olema globals(), muidu Airflow seda ei leia. DAG'ile tuleb ka öelda:
    • kuidas seda kutsuda orders see nimi ilmub hiljem veebiliideses,
    • et see hakkab töötama alates 8. juuli keskööst,
    • ja kĂ€ivitama ta peaks umbes iga 6 tunni jĂ€rel (karmimatele on siin 'timedelta()' asemel lubatud -string, kergematele - vĂ€ljend nagu cron@daily 0 0 0/6 ? * * *workflow() teeb peamise töö, aga mitte praegu. Praegu viskame lihtsalt meie konteksti logisse.);
  • NĂŒĂŒd lihtne maagia ĂŒlesannete loomisel: jalutame meie allikate kaudu;
  • initsialiseerime
    • , mis tĂ€idab meie anomaalia,
    • . Ärge unustage mÀÀrata ĂŒlesande ainulaadne (DAG'i ulatuses) nimi ja seondada see DAG'iga. Lipp PythonOperatorprovide_context NĂŒĂŒd lihtne maagia ĂŒlesannete loomisel:toob omakorda funktsiooni lisaargumente, mille me hoolikalt kogume, kasutades **context Sellega on praegu kĂ”ik. Mida me saime: uus DAG veebiliideses,.

poolteist sada ĂŒlesannet, mis töötavad paralleelselt (kui Airflow, Celery ja serverite seadistused seda lubavad).

  • Noh, peaaegu saime.
  • Kes seab sĂ”ltuvused paika?

Selle asja lihtsustamiseks lĂŒlitasin selle

Apache Airflow: lihtsustame ETL-i
kÀsitlemiseks

kĂ”ikidesse sĂ”lmedesse. docker-compose.yml NĂŒĂŒd siis lĂ€ks lahti: requirements.txt Hallid ruudud - ĂŒlesande instantsid, mida planeerija on töödelnud.

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

Apache Airflow: lihtsustame ETL-i

Rohelised, nagu mÔistetav, - edukalt lÔpule viidud. Punased - mitte nii edukalt.

Muide, meie tootel ei ole mingit kausta,

Apache Airflow: lihtsustame ETL-i

, mis sĂŒsynkroniseerib masinate vahel - kĂ”ik DAG'id asuvad

meie Gitlab'is, ja Gitlab CI paigutab uuendused masinatesse, kui toimub liitmine . /dagsKui töötajad töötlevad meie tĂŒhiseid ĂŒlesandeid, tuletame meelde veel ĂŒhte tööriista, mis meile midagi nĂ€idata vĂ”ib - Flower. git Esimene leht, kus on kokkuvĂ”tvad andmed sĂ”lmedest - töötajatest: master.

Natuke Flowerist

KĂ”ige jaotatum leht ĂŒlesannetest, mis on tööle saadetud:

Esimene leht, mis sisaldab kokkuvÔtlikku teavet töötajate sÔlmpunktide kohta:

Apache Airflow: lihtsustame ETL-i

Rikkaim leht ĂŒlesannetega, mis on saadetud töötlemiseks:

Apache Airflow: lihtsustame ETL-i

KÔige igavam leht meie maakleri seisundi kohta:

Apache Airflow: lihtsustame ETL-i

KĂ”ige vĂ€rvilisem leht — ĂŒlesannete seisundi graafikute ja nende tĂ€itmise ajaga:

Apache Airflow: lihtsustame ETL-i

Laadime alla alalaadimise

Nii et kĂ”ik ĂŒlesanded on töötanud, saame haavatud Ă€ra viia.

Apache Airflow: lihtsustame ETL-i

Haavatuid oli ĂŒllatavalt palju — erinevatel pĂ”hjustel. Kui Airflowd kasutatakse Ă”igesti, siis need ruudud nĂ€itavad, et andmed ei ole kindlasti kohale jĂ”udnud.

Peame vaatama logisid ja taaskĂ€ivitama nurjunud ĂŒlesande nĂ€idiseid.

Klikkides igal ruudul, nÀeme meie kÀsutuses olevaid toiminguid:

Apache Airflow: lihtsustame ETL-i

Saame vĂ”tta ja teha nurjunule Clear. See tĂ€hendab, et unustame, et seal on mingi probleem ja sama ĂŒlesande nĂ€idis suundub planeerijale.

Apache Airflow: lihtsustame ETL-i

Selge, et ei ole vĂ€ga inimsĂ”bralik kĂ”iki punaseid ruute hiirega nii kĂ€sitleda — seda me Airflowlt ei oota. Loomulikult on meil massihĂ€vitusrelv: Browse/Task Instances

Apache Airflow: lihtsustame ETL-i

Valime kÔik korraga ja nullime, valides Ôige punkti:

Apache Airflow: lihtsustame ETL-i

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

Apache Airflow: lihtsustame ETL-i

Ühendused, hookid ja muud muutujad

On paras 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("""Lugupeetud, 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, meil on {{ dag.dag_id }} kukkunud
        """),
    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 kÔik on kunagi teinud raportite uuendamist? See on jÀlle see: on nimekiri allikatest, kust andmeid vÔtta; on nimekiri, kuhu need panna; Àrge unustage mÀrku anda, kui kÔik juhtus vÔi lÀks rikki (noh, see ei puuduta meid, ei).

Vaatame taas faili ja uurime uusi arusaamatuid asju:

  • from commons.operators import TelegramBotSendMessage — ei ole midagi, mis takistaks meil oma operaatorite loomist, mida me tegime, luues vĂ€ikese mĂ€hise sĂ”numite saatmiseks Unblockitud. (Sellest operaatorist rÀÀgime ka hiljem);
  • default_args={} — DAG vĂ”ib jagada samu argumendi kĂ”igile oma operaatoritele;
  • to='{{ var.value.all_the_kings_men }}' — vĂ€li to meil ei ole seda kĂ”vasti mÀÀratletud, vaid see genereeritakse dĂŒnaamiliselt Jinja ja e-mailide nimekirja muutuja abil, mille ma hoolikalt asetasin Admin/Variables;
  • trigger_rule=TriggerRule.ALL_SUCCESS — operaatori kĂ€ivitamise tingimus. Meie puhul saadetakse kiri juhtidele ainult siis, kui kĂ”ik sĂ”ltuvused töötasid edukalt;
  • tg_bot_conn_id='tg_main' — argumendid conn_id vĂ”tavad endasse ĂŒhenduste identifikaatorid, mille me loome Admin/Connections;
  • trigger_rule=TriggerRule.ONE_FAILED — Telegrami sĂ”numid saadetakse ainult, kui on katkestatud ĂŒlesandeid;
  • task_concurrency=1 — keelame sama ĂŒlesande mitme ĂŒlesande eksemplari samaaegse kĂ€ivitamise. Vastasel juhul saame mitu korraga VerticaOperator (mis vaatavad ĂŒhte tabelit);
  • report_update >> [email, tg] — kĂ”ik VerticaOperator kohtuvad kirja ja sĂ”numi saatmisel, just nii:
    Apache Airflow: lihtsustame ETL-i

    Kuid kuna teavitajate operaatoritel on erinevad kĂ€ivitamistingimused, töötab ainult ĂŒks. Puustruktuuris nĂ€eb see vĂ€lja vĂ€hem selge:
    Apache Airflow: lihtsustame ETL-i

RÀÀgin paar sĂ”na makrode ja nende sĂ”prade — muutujad.

Makrod on Jinja-plekistrid, mis vÔivad esitada 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 }} lahti muudetakse konteksti muutuja sisuks execution_date Port 2222 YYYY-MM-DD: 2020-07-14. KĂ”ige parem on see, et konteksti muutujaid seotakse kindlalt konkreetse ĂŒlesande eksemplariga (ruudukesega Puustruktuuris), ja uuesti kĂ€ivitamisel avanevad plekistrid samadele vÀÀrtustele.

MÀÀratud vÀÀrtusi saab vaadata nuppuga Rendered igal ĂŒlesande eksemplaril. NĂ€iteks on see ĂŒlesandel, mis saadab kirja:

Apache Airflow: lihtsustame ETL-i

Ja see ĂŒlesandel, mis saadab sĂ”numi:

Apache Airflow: lihtsustame ETL-i

TĂ€ielik nimekiri sisseehitatud makrodest viimase saadaval oleva versiooni jaoks on siin: Macros Reference

Veelgi enam, pluginade abil saame kuulutada oma makrosid, kuid see on juba hoopis teine lugu.

Lisaks eelnevalt mÀÀratud asjadele saame kasutada oma muutujaid (muutsin neid ĂŒlaltoodud koodis). Loome Admin/Variables paar asja:

Apache Airflow: lihtsustame ETL-i

KĂŒll, nĂŒĂŒd saab hakata kasutama:

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

MÔisted vÔivad olla skalaarsed vÔi JSON. JSON-i puhul:

bot_config

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

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

KĂ”nestan vaid ĂŒhe sĂ”na ja nĂ€itan ĂŒhte ekraanipilti seoses ĂŒhendustega. Siin on kĂ”ik elementaarne: lehe peal Admin/Connections loome ĂŒhenduse, paneme sinna oma kasutajanimed/paroolid ja veel spetsiifilisemad seaded. Nii:

Apache Airflow: lihtsustame ETL-i

Paroolid saab krĂŒpteerida (rohkem kui vaikimisi juhul) vĂ”i pole mÀÀratud ĂŒhenduse tĂŒĂŒpi (nagu tegin tg_main) — asi on selles, et tĂŒĂŒbid on Airflow mudelites sisse ehitatud ja ilma lĂ€htekoodidesse sekkumata neid ei saa muuta (kui midagi on vale, siis palun parandage mind), kuid krĂŒpteerida nime jĂ€rgi pole meil midagi takistust.

Ja lisaks vĂ”ime luua mitu sama nimega ĂŒhendust: sel juhul meetod BaseHook.get_connection(), mis toob meile ĂŒhendused nime jĂ€rgi, annab juhusliku ĂŒhe mitme sarnase seast (oleks loogilisem teha Round Robin, kuid jĂ€tame selle Airflow arendajate mureks).

Muutujad ja ĂŒhendused on kahtlemata suurepĂ€rased tööriistad, kuid oluline on mitte kaotada tasakaalu: millised teie voogude osad sĂ€ilitad konkreetselt koodis ning millised usaldad Airflow’d. Ühelt poolt on mugav kiirelt muuta vÀÀrtust, nĂ€iteks postkasti, lĂ€bi UI. Teiselt poolt on see ikka naasmine hiireklikkide juurde, millest tahtsime (mina) lahti saada.

Ühendustega töötamine on ĂŒks ĂŒlesanne hookide paralleelseks tĂ€itmiseks.. Üldiselt on Airflow hook'id ĂŒhenduspunktid kolmandate teenuste ja raamatukogudega. NĂ€iteks JiraHook avab meile kliendi suhtlemiseks Jira'ga (saame ĂŒlesandeid liigutada edasi-tagasi), ja SambaHook vĂ”imaldab me pushida kohaliku faili smb-punkti.

TĂ€itame kohandatud operaatorit

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

Kood commons/operators.py oma operaatori kohta:

importida 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'Saadan "{self.message}" chati {self.chat_id}')
        self.client.send_message(chat_id=self.chat_id,
                                 message=self.message)

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

  • Oleme pĂ€rinud BaseOperator, mis rakendab palju Airflow-spetsiifilisi asju (vaadake kunagi, kui on aega)
  • Oleme kuulutanud vĂ€lja vĂ€ljad template_fields, kus Jinja otsib makrosid töötlemiseks.
  • Oleme korraldanud Ă”iged argumendid __init__(), seadnud vaikimisi vÀÀrtused sinna, kus on vajalik.
  • Ka eelneva initsialiseerimist ei unustanud.
  • Oleme avanud vastava hook'i TelegramBotHook, saanud sellelt kliendi objekti.
  • Olemegi ĂŒle kirjutanud (override) meetodi BaseOperator.execute(), mida Airflow kutsub vĂ€lja, kui on aeg operaatorit tööle panna — just seal me teeme pĂ”hitegevuse, unustamata logida. (Me logime, muide, otse stdout ja stderr — Airflow teeb kĂ”ik pĂŒĂŒdmiseks ja pakib ilusasti kokku, paneb Ă”igesse kohta.)

Vaatame, mis meil on commons/hooks.py. Faili esimene osa, koos hook'iga:

importida 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 ei tea isegi, mida siit selgitada, lihtsalt mÀrgin olulised punktid:

  • PĂ€rime, mĂ”tleme argumentidele — enamasti on neid ĂŒks: conn_id;
  • Üle kirjutame standardsed meetodid: ma piirdusin get_conn(), kus ma saan ĂŒhenduse parameetrid nime jĂ€rgi ja lihtsalt tĂ”stan sektsiooni. extra (see code for JSON), where I put the Telegram bot token according to my own instructions: {"bot_token": "YOuRAwEsomeBOtToKen"}.
  • Creating an instance of our TelegramBot, giving it a specific token.

That's it. You can get the client from the hook using TelegramBotHook().client vÔi TelegramBotHook().get_conn().

And the second part of the file, where I wrap the Telegram REST API to avoid carrying the same python-telegram-bot for just one method 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))

The correct approach is to put all of this in: TelegramBotSendMessage, TelegramBotHook, TelegramBot — a plugin, store it in a public repository, and release it as Open Source.

While we were studying all this, our report updates successfully piled up and sent me a message about an error in the channel. I will go check what went wrong again...

Apache Airflow: lihtsustame ETL-i
Something broke in our DAG! Wasn't this what we were waiting for? Exactly!

Kas sa teed seda?

Do you feel like I've missed something? I promised to transfer data from SQL Server to Vertica, and here I went off-topic, how careless!

This misdeed was intentional; I simply had to explain some terminology to you. Now we can move on.

Our plan was as follows:

  1. Create a DAG
  2. Generate tasks
  3. See how everything looks nice
  4. Assign session numbers to uploads
  5. Fetch data from SQL Server
  6. Store data in Vertica
  7. Compile statistics

So, to run all this, I made a small addition to our 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 tÔstame:

  • Vertica kui host dwh kĂ”ige vaikimisi seadistustega,
  • kolm SQL Serveri eksemplari,
  • tĂ€iendame andmebaase viimasel ajal mĂ”ne andmega (Ă€rge mingil juhul vaadake sisse mssql_init.py!)

KÀivitame kogu selle hea kraami veidi keerukama kÀsuga kui eelmisel korral:

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

Mis meie imetore genereerija tootis, on vÔimalik, kasutades punkti Data Profiling/Ad Hoc Query:

Apache Airflow: lihtsustame ETL-i
Peamine, et seda analĂŒĂŒtikutele ei nĂ€idata

Detailid ETL-seanssides ma ei peatuks, seal on kĂ”ik triviaalne: loome andmebaasi, sinna tabeli, katame kĂ”ik konteksti juhiga, ja nĂŒĂŒd teeme nii: with Session(task_name) as session: print('Load', session.id, 'started')# Load workflow ...session.successful = True session.loaded_rows = 15

session.py

session.py

from sys import stderr

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

    NĂ€ide:
        with Session(task_name) as session:
            print(session.id)
            session.successful = True
            session.loaded_rows = 15
            session.comment = 'Hea 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 = """
            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, 'avatud')
        return self

    def close(self):
        if not self._id:
            raise SessionClosedError('Seanss 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,
              ', Laaditud: ', self.loaded_rows,
              ', kommentaar:', self.comment)

class SessionError(Exception):
    pass

class SessionClosedError(SessionError):
    pass

KÀtte on jÔudnud aeg andmed kokku koguda meie poolteise saja tabeli hulgast. Teeme seda vÀga lihtsalt ridadega:

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. HĂŒĂŒdme abil saame Airflow'st pymssql-ĂŒhenduse
  2. KĂŒsitlusse lisame kuupĂ€eva piirangu - selle edastab meile mallitegija.
  3. KĂ€ivitame meie pĂ€ringu pandas, mis tĂ”mbab meile DataFrame — see tuleb meile hiljem kasuks.

Kasutame asendust {dt} kĂŒsi parameetri asemel %s mitte sellepĂ€rast, et ma olen kuri Buratino, vaid pigem sellepĂ€rast, et pandas ei suuda hakkama saada pymssql ja paneb viimasele params: Loend, kuigi see vĂ€ga tahab tulp.
Pange tÀhele, et arendaja pymssql otsustas teda enam mitte toetada, ja on aeg kolida pyodbc.

Vaatame, millega Airflow tÀitis meie funktsioonide argumendid:

Apache Airflow: lihtsustame ETL-i

Kui andmeid ei olnud, siis pole mĂ”tet jĂ€tkata. Aga ka selle ĂŒle lugemine, et laadimine Ă”nnestus, on kummaline. Kuid see pole viga. A-a-a, mida teha?! Siin on, mida:

if df.empty:
    raise AirflowSkipException('Ridasid ei ole laadimiseks')

AirflowSkipException ĂŒtleb Airflow, et viga ei ole, ja ĂŒlesanne jÀÀb vahele. Liideses ei ole roheline ja ei punane ruut, vaid roosa.

Lisame meie andmetele mÔned veerud:

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

TĂ€pselt:

  • DB, kust me tellimused vĂ”tsime,
  • Meie laadimisessiooni identifikaator (see on erinev iga ĂŒlesande jaoks),
  • Allika ja tellimuse identifikaatori hash — et lĂ”pp-andmebaasis (kus kĂ”ik valatakse ĂŒhte tabelisse) oleks meil ainulaadne tellimuse identifikaator.

JÀÀnud on eelviimane samm: laadida kĂ”ik Vertica'sse. Ja, nagu ei oleks uskumatu, on ĂŒks efektiivseid viise teha seda — lĂ€bi CSV!

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

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

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

cursor = target_conn.cursor()
cursor.copy(copy_stmt, buffer)
  1. Teeme spetsiaalse vastuvÔtja StringIO.
  2. pandas mis lahendab meie DataFrame vÀlise lisandina CSV-rid.
  3. Avame ĂŒhenduse meie lemmik Vertica hook'iga.
  4. Ja nĂŒĂŒd saame copy() saata meie andmed otse Vertica'sse!

Draiverist vĂ”tame, kui palju ridu laaditi, ja ĂŒtleme sessiooni juhile, et kĂ”ik on OK:

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

See on kÔik.

Produksioonis loome sihttableti kÀsitsi. Siin lubasin endale vÀikese automaatika:

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 kasutan VerticaOperator() loomiseks andmebaasi skeemi ja tabeli (kui need veel puuduvad, loomulikult). 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Ôtte

— No nii, — ĂŒtles hiirepoeg, — ei ole tĂ”si, et nĂŒĂŒd
Kas sa oled veendunud, et metsa kÔige hirmsam loom olen mina?

Julia Donaldson, „Gruffalo“

MÔtle, kui me kolleegidega korraldaksime konkursi: kes suudab kiiremini luua ja kÀivitada ETL-protsessi nullist: nemad oma SSIS-i ja hiirega ja mina Airflow'iga... Ja siis vÔrdleksime veel hooldamise mugavust... Oh, ma arvan, et oled nÔus, et ma möödun neist igas osas!

Kui rÀÀkida natuke tĂ”sisemalt, siis Apache Airflow – tĂ€nu protsesside kirjeldamisele programmikoodina – tegi mu töö mĂ€rksa mugavamaks ja meeldivamaks.

Selle piiramatud laienemisvĂ”imalused: nii pluginate osas kui ka skaleerimisvĂ”imes — annavad vĂ”imaluse rakendada Airflow’d praktiliselt igas valdkonnas: olgu see siis andmete kogumise, ettevalmistamise ja töötlemise tĂ€issĂŒkkel vĂ”i isegi rakettide kĂ€ivitamine (loomulikult Marsile).

Viimane, teaduslik-informatiivne osa

Kohad, mille me teie jaoks kokku korjasime

  • start_date. Jah, see on juba kohalik meem. Peamise DAG-argumenti 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 pole mingeid probleeme.

    Sellega on seotud veel ĂŒks tĂ€itmisviga: Task is missing the start_date parameter, mis enamikul juhtudest ĂŒtleb, et unustasid DAG-i operaatoriga siduda.

  • KĂ”ik ĂŒhel masinal. Jah, ja andmebaasid (nii Airflow enda kui ka meie ĂŒmbruse), ja veebiserver, ja ajastaja, ja töötajad. Ja see isegi toimis. Aga aja jooksul kasvas teenuste ĂŒlesannete arv, ja kui PostgreSQL hakkas vastama indeksilt 20 ms asemel 5 ms, vĂ”tsime selle ja viisime minema.
  • LocalExecutor. Jah, me kasutame seda endiselt ja oleme juba ÀÀre peal. LocalExecutori jaoks on meil siiani piisavalt, kuid nĂŒĂŒd on aeg vĂ€hemalt ĂŒhe töötaja vĂ”rra laiendada, ja tuleb pingutada, et ĂŒle minna CeleryExecutori peale. Ja arvestades, et sellega saab töötada ka ĂŒhel masinal, ei takista miski Celery kasutamist isegi mitte serveris, mis "loomulikult ei tule kunagi tootmisse, tĂ”otame!"
  • Kasutamata sisekasutuse:
    • Connections teenuste autentimisandmete haldamiseks,
    • SLA Misses ĂŒlesannete jaoks, mis ei töötanud Ă”igeaegselt,
    • XCom metatekstide vahetamiseks (ma ĂŒtlesin metateavet!) DAG-i ĂŒlesannete vahel.
  • Posti kuritarvitamine. Mis siin ikka öelda? Olin seadnud teavitused kĂ”ikide korduvate ebaĂ”nnestunud ĂŒlesannete jaoks. NĂŒĂŒd on minu töö Gmailis ĂŒle 90 000 e-kirja Airflow'lt ning postiteenuse veebivĂ€ljaanne keeldub kustutamast rohkem kui 100 korraga.

Rohkem peidetud kivisid: Apache Airflowi Probleemid

Veel suurema automatiseerimise vahendid

Selleks, et saaksime veel rohkem mÔelda, mitte ainult kÀtega töötada, on Airflow meie jaoks ette valmistanud jÀrgmised vÔimalused:

  • REST API — tal on endiselt Eksperimentaalse staatuse, mis ei takista selle toimimist. Selle abil on vĂ”imalik mitte ainult saada teavet DAG-ide ja ĂŒlesannete kohta, vaid ka peatada/taasalustada DAG-i, luua DAG Run vĂ”i baas.
  • cli-runtime'is ja kubectl'is — kĂ€surealt on saadaval palju tööriistu, mis ei ole lihtsalt ebamugavad veebiliidese kaudu kasutada, vaid pole seal isegi saadaval. NĂ€iteks:
    • backfill vajalik, et kĂ€ivitada ĂŒlesannete instantside uuesti kĂ€ivitamine.
      NĂ€iteks tulid analĂŒĂŒtikud ja ĂŒtlesid: „Teie, seltsimees, andmetes on jama 1. kuni 13. jaanuarini! Parandage!“. Ja siis sa ĂŒtled:
      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 kĂ”rvale. Veelgi enam, seda saab kĂ€ivitada lĂ€bi LocalExecutor, isegi kui sul on Celery klastri.
    • Umbes sama asja teeb test, ainult et see ei kirjuta andmebaasi.
    • connections vĂ”imaldab massiliselt luua ĂŒhendusi shellist.
  • Python API on ĂŒsna keeruline suhtlemisviis, mis on mĂ”eldud pistikprogrammidena, mitte selle kĂ€sitsi mudimiseks. Aga kes meid takistab minemast /home/airflow/dags, kĂ€ivitama ipython ja hakkama siin mĂ€ngima? NĂ€iteks vĂ”ib kĂ”iki ĂŒhendusi eksportida jĂ€rgmise 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)
  • Ühendus Airflow'i metadatuuri andmebaasi. Ma ei soovita sinna kirjutada, kuid ĂŒlesannete olekute saamine erinevate spetsiifiliste metrikate jaoks on oluliselt kiirem ja lihtsam kui lĂ€bi ĂŒhegi API.

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

    Ole ettevaatlik, SQL!

    VIIMASTE_EXECUTIONIDEGA AS (
    VALI
        ĂŒlesande_id,
        dag_id,
        tÀitmise_aeg,
        olek,
            rida_numbrina()
            ÜLE (
                JAOTUSEKS ĂŒlesande_id, dag_id
                KORRALDA tÀitmise_aeg LANGUS) NII rn
    FROM public.task_instance
    KUS
        tĂ€itmise_aeg > nĂŒĂŒd() - INTERVALL '2' PÄEVA
    ),
    ebaÔnnestunud AS (
        VALI
            ĂŒlesande_id,
            dag_id,
            tÀitmise_aeg,
            olek,
            JUHTUM KUI rn = rida_numbrina() ÜLE (
                JAOTUSEKS ĂŒlesande_id, dag_id
                KORRALDA tÀitmise_aeg LANGUS)
                     SIIS TÕENE LÕPP JN kui ebaĂ”nnestunud
        KUS
            olek IN ('ebaÔnnestunud', 'ootab_kordamist')
    )
    VALI
        ĂŒlesande_id,
        dag_id,
        loe(ebaÔnnestunud) AS ebaÔnnestunud,
        loe(JUHTUM KUI ebaÔnnestunud
            JA olek = 'ebaĂ”nnestunud' SIIS 1 LÕPP) AS ebaĂ”nnestunud,
        loe(JUHTUM KUI ebaÔnnestunud
            JA olek = 'ootab_kordamist' SIIS 1 LÕPP) AS ootab_kordamist
    FROM ebaÔnnestunud
    GRUPEERI
        ĂŒlesande_id,
        dag_id
    OLLES
        loe(ebaÔnnestunud) > 0

Viidatud lingid

Ja loomulikult esimesed kĂŒmme linki Google'i otsingust, mis viivad minu Airflow kaustadesse.

Ja lingid, mis on artiklis kasutatud:

Allikas: habr.com

Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid đŸ”„ Osta usaldusvÀÀrne hostimine veebilehtede jaoks DDoS-i kaitsega, VPS VDS serverid | ProHoster