Apache Airflow: e bëjmë ETL më të lehtë

PĂ«rshĂ«ndetje, unĂ« jam Dmitriy Logvinenko — Inxhinier i tĂ« DhĂ«nave nĂ« departamentin e analizĂ«s sĂ« grupit tĂ« kompanive «VezĂ«t».

Do t'ju tregoj pĂ«r njĂ« mjet tĂ« shkĂ«lqyer pĂ«r zhvillimin e proceseve ETL — Apache Airflow. Por Airflow Ă«shtĂ« aq universal dhe shumĂ«fishtĂ«, saqĂ« ia vlen ta shqyrtoni edhe nĂ«se nuk mereni me rrjedhat e tĂ« dhĂ«nave, por keni nevojĂ« tĂ« filloni nganjĂ«herĂ« ndonjĂ« proces dhe tĂ« ndihmoni nĂ« mbajtjen e tyre nĂ«n kontroll.

Po, nuk do të flas vetëm, por do të tregoj dhe: programi ka shumë kode, shkëmbime dhe rekomandime.

Apache Airflow: e bëjmë ETL më të lehtë
ÇfarĂ« zakonisht shihni kur kĂ«rkoni fjalĂ«n Airflow / Wikimedia Commons

Përmbajtja

Hyrje

Apache Airflow — Ă«shtĂ« ashtu si Django:

  • shkruar nĂ« Python,
  • ka njĂ« admin tĂ« shkĂ«lqyer,
  • Ă«shtĂ« pa kufij pĂ«r zgjerim,

— vetĂ«m mĂ« tĂ« mirat, dhe Ă«shtĂ« krijuar pĂ«r qĂ«llime krejtĂ«sisht tĂ« tjera, pĂ«rkatĂ«sisht (siç Ă«shtĂ« shkruar deri nĂ« kaptinĂ«):

  • nisjen dhe monitorimin e detyrave nĂ« numĂ«r tĂ« pakufizuar makinash (sa do t'ju lejojĂ« Celery/Kubernetes dhe ndĂ«rgjegjja juaj)
  • me gjenerimin dinamik tĂ« flukseve tĂ« punĂ«s nga njĂ« kod shumĂ« i lehtĂ« pĂ«r t'u shkruar dhe pĂ«r t'u kuptuar nĂ« Python
  • dhe me mundĂ«sinĂ« pĂ«r tĂ« lidhur çdo bazĂ« tĂ« tĂ« dhĂ«nave dhe API me ndihmĂ«n e komponenteve tĂ« gatshme ose pluginĂ«ve tĂ« vetĂ«-krijuar (gjĂ« qĂ« bĂ«het jashtĂ«zakonisht e lehtĂ«).

Ne përdorim Apache Airflow në këtë mënyrë:

  • mbledhim tĂ« dhĂ«na nga burime tĂ« ndryshme (shumĂ« instanca SQL Server dhe PostgreSQL, API tĂ« ndryshme me metrika tĂ« aplikacioneve, madje edhe 1C) nĂ« DWH dhe ODS (kĂ«tĂ« e bĂ«jmĂ« me Vertica dhe Clickhouse).
  • si njĂ« avanguardĂ« cron, i cili ekzekuton procese konsolidimi tĂ« tĂ« dhĂ«nave nĂ« ODS, si dhe monitoron shĂ«rbimin e tyre.

Derisa përpara pak kohësh ne e përmbushëm nevojën tonë me një server të vogël me 32 bërthama dhe 50 GB RAM. Në Airflow funksionon:

  • mĂ« 200 DAG-Ă« (nĂ« thelb flukse pune, nĂ« tĂ« cilat kemi mbushur detyrat),
  • nĂ« secilin nĂ« mesatarisht 70 detyra,
  • kjo gjĂ« ekzekutohet (po ashtu nĂ« mesatarisht) njĂ« herĂ« nĂ« orĂ«.

Dhe pĂ«r mĂ«nyrĂ«n se si ne u zgjeruam, do tĂ« shkruaj mĂ« poshtĂ«, por tani le tĂ« pĂ«rcaktojmĂ« ĂŒber-detyrĂ«n qĂ« do tĂ« zgjidhim:

Ka janĂ« tri servera SQL burimorĂ«, tĂ« cilĂ«t kanĂ« nga 50 baza tĂ« dhĂ«nash — instanca tĂ« njĂ« projekti, pĂ«rkatĂ«sisht struktura e tyre Ă«shtĂ« e njĂ«jtĂ« (nĂ« shumicĂ«n e rasteve, mwa-ha-ha), dhe pĂ«r pasojĂ«, çdo njĂ«ri ka njĂ« tabelĂ« Orders (fatmirĂ«sisht, njĂ« tabelĂ« me njĂ« emĂ«r tĂ« tillĂ« mund tĂ« futet nĂ« çdo biznes). Ne marrim tĂ« dhĂ«nat duke shtuar fushat ndihmĂ«se (serveri burim, baza burimore, identifikuesi i detyrĂ«s ETL) dhe naivisht do t'i hedhim ato nĂ«, tĂ« themi, Vertica.

Fillojmë!

Pjesa kryesore, praktike (edhe pak teorike)

Pse na duhet (dhe juve)

Kur pemët ishin të mëdha dhe unë isha një i thjeshtë SQL-çkuar në një shitje të madhe ruse, ne hidheshim në proceset ETL aka rrjedhat e të dhënave me dy mjete të disponueshme për ne:

  • Informatica Power Center — njĂ« sistem ekstremisht kompleks, jashtĂ«zakonisht produktiv, me pajisjet e veta dhe versionimin e saj. Kam pĂ«rdorur ndoshta 1% tĂ« mundĂ«sive tĂ« saj. Pse? Epo, sĂ« pari, ky ndĂ«rfaqe ishte diku nga vitet '00 dhe na e shkaktonte presion psikologjik. SĂ« dyti, kjo gjĂ« Ă«shtĂ« e mbyllur nĂ« procese tejet komplekse, riblerje intensive tĂ« komponenteve dhe shumĂ« karakteristika tĂ« rĂ«ndĂ«sishme pĂ«r ndĂ«rmarrjet. PĂ«r atĂ« çmim qĂ« kushton, si fluturimi i njĂ« Airbus A380/njĂ« vit, do tĂ« heshtim.

    Kujdes, screenshot mund t'i bëjë të rinjve nën 30 vjeç të ndihen disi keq

    Apache Airflow: e bëjmë ETL më të lehtë

  • SQL Server Integration Server — ne kemi pĂ«rdorur kĂ«tĂ« shok nĂ« proceset tona interne. Por nĂ« realitet: ne tashmĂ« pĂ«rdorim SQL Server, dhe do tĂ« ishte disi e paarsyeshme tĂ« mos i shfrytĂ«zonim mjetet ETL tĂ« tij. Gjithçka nĂ« tĂ« Ă«shtĂ« mirĂ«: dhe ndĂ«rfaqja Ă«shtĂ« e bukur, si dhe raportet e ekzekutimit... Por nuk pĂ«r kĂ«tĂ« e duam produktin software, oh jo pĂ«r kĂ«tĂ«. TĂ« versiononi atĂ« dtsx (i cili Ă«shtĂ« njĂ« XML me node tĂ« pĂ«rzier gjatĂ« ruajtjes) mund ta bĂ«jmĂ«, por çfarĂ« do tĂ« thotĂ«? TĂ« bĂ«jmĂ« njĂ« paketĂ« detyrash qĂ« do tĂ« transferonte njĂ«qind tabela nga njĂ« server nĂ« tjetrin? Po çfarĂ« njĂ«qind, do t'ju bie gishti tregues pas njĂ«zet copash, duke klikuar mbi butonin e mausit. Por duke parĂ«, sigurisht duket mĂ« modern:

    Apache Airflow: e bëjmë ETL më të lehtë

Kemi kërkuar pa dyshim zgjidhje. Punë madje gjysmë arritëm të krijojmë një gjenerator të vetë-ndërtuar për paketat SSIS...

... dhe pastaj më gjeti puna e re. Dhe aty më preku Apache Airflow.

Kur mĂ«sova se pĂ«rshkrimet e proceseve ETL janĂ« thjesht kod Python, pothuajse kĂ«rcen nga gĂ«zimi. KĂ«shtu qĂ« rrjedhat e tĂ« dhĂ«nave u nĂ«nshtruan versionimit dhe difes, dhe grumbullimi i tabelave me njĂ« strukturĂ« unike nga njĂ«qind baza tĂ« dhĂ«nash nĂ« njĂ« target u bĂ« njĂ« punĂ« e kodit Python nĂ« njĂ« ekran 13” me njĂ« tĂ« dhjetĂ«.

Po mblidhni një klaster

Le të mos organizojmë një kopësht të fëmijëve, dhe të mos flasim për gjëra të heltësisht të qarta, siç është instalimi i Airflow, baza e të dhënave që keni zgjedhur, Celery dhe punë të tjera të përshkruara në dokumente.

Që të mund të fillojmë menjëherë me eksperimentet, kam skicuar docker-compose.yml në të cilin:

  • Le tĂ« ngrisim vetĂ« Airflow: Scheduler, Webserver. Atje do tĂ« funksionojĂ« gjithashtu Flower pĂ«r monitorimin e detyrave tĂ« Celery (sepse Ă«shtĂ« shtypur tashmĂ« nĂ« apache/airflow:1.10.10-python3.7, dhe ne nuk kemi asnjĂ« problem me kĂ«tĂ«);
  • PostgreSQL, nĂ« tĂ« cilin Airflow do tĂ« shkruajĂ« informacionin e saj operativ (tĂ« dhĂ«nat e planifikuesit, statistikat e ekzekutimit etj.), ndĂ«rsa Celery do tĂ« shĂ«nojĂ« detyrat e pĂ«rfunduara;
  • Redis, i cili do tĂ« veprojĂ« si njĂ« broker detyrash pĂ«r Celery;
  • Celery worker, i cili do tĂ« merret me ekzekutimin e drejtpĂ«rdrejtĂ« tĂ« detyrave.
  • NĂ« dosjen ./dags ne do tĂ« vendosim skedarĂ«t tanĂ« me pĂ«rshkrimin e dags. Ato do tĂ« kapen nĂ« flakĂ«, prandaj nuk Ă«shtĂ« e nevojshme tĂ« tĂ«rhiqni tĂ« gjithĂ« stekun pas çdo shkurti.

Disa herë kodi në shembuj është dhënë jo plotësisht (për të mos e mbingarkuar tekstin), dhe diku ai modifikohet gjatë procesit. Shembujt e plotë funksionues të kodit mund të shihen në repositorium. 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 si Celery broker
  broker:
    image: redis:6.0.5-alpine

  # DB për metadata e Airflow
  airflow-db:
    image: postgres:10.13-alpine

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

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

  # Kontejner kryesor me Webserver, Scheduler, Celery Flower të Airflow
  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
      # Webserver i Airflow
      - 8080:8080

  # Punëtori Celery, do të shkallëzohet duke përdorur `--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

Shënime:

  • NĂ« ndĂ«rtimin e kompozitĂ«s kam mbĂ«shtetur shumĂ« nĂ« imazhin e njohur puckel/docker-airflow – patjetĂ«r shikoni. Ndoshta nĂ« jetĂ«n tuaj nuk keni nevojĂ« pĂ«r asgjĂ« tjetĂ«r.
  • TĂ« gjitha konfigurimet e Airflow janĂ« tĂ« disponueshme jo vetĂ«m pĂ«rmes airflow.cfg, por edhe pĂ«rmes variablave tĂ« ambientit (falĂ« zhvilluesve), tĂ« cilat i pĂ«rdora keq.
  • Natyrisht, ai nuk Ă«shtĂ« gati pĂ«r prodhim: qĂ«llimisht nuk kam vendosur heartbeats nĂ« kontenierĂ«, nuk u shqetĂ«sova pĂ«r sigurinĂ«. MegjithatĂ«, bĂ«ra njĂ« minimum tĂ« pĂ«rshtatshĂ«m pĂ«r eksperimentet tona.
  • VĂ« re se:
    • Folderi me DAG-at duhet tĂ« jetĂ« i aksesueshĂ«m si pĂ«r planifikuesin ashtu edhe pĂ«r punĂ«torĂ«t.
    • E njĂ«jta gjĂ« vlen edhe pĂ«r tĂ« gjitha bibliotekat e palĂ«ve tĂ« treta – ato duhet tĂ« jenĂ« tĂ« gjitha tĂ« instaluara nĂ« makinat me planifikuesin dhe punĂ«torĂ«t.

Tani thjesht:

$ docker-compose up --scale worker=3

Pas ngritjes, mund të shikoni ndërfaqet e internetit:

Koncepte themelore

Nëse nuk keni kuptuar asgjë nga këto «DAG», ja një fjalor i shkurtër:

  • Planifikuesi – ŰŻŰłŰȘßgĆ· (djali mĂ« i rĂ«ndĂ«sishĂ«m) nĂ« Airflow, qĂ« kontrollon qĂ« robotĂ«t tĂ« punojnĂ«, jo njerĂ«zit: monitoron orarin, pĂ«rditĂ«son DAG-at, nis detyrat.

    NĂ« versionet mĂ« tĂ« vjetra, ai kishte probleme me memorinĂ« (jo, jo amnezia, por rrjedhje) dhe atje kishte mbetur njĂ« parametĂ«r legacy nĂ« konfigurime. run_duration — intervali i rivendosjes sĂ« tij. Por tani gjithçka Ă«shtĂ« nĂ« rregull.

  • DAG (edhe i njohur si «dag») — «grafik i orientuar jo-ciklik», por njĂ« pĂ«rcaktim i tillĂ« shumĂ« pak njerĂ«zve do t'u thotĂ« diçka, dhe nĂ« thelb, Ă«shtĂ« njĂ« kontejner pĂ«r detyrat qĂ« ndĂ«rveprojnĂ« me njĂ«ra-tjetrĂ«n (shih mĂ« poshtĂ«) ose njĂ« analog i PaketĂ«s nĂ« SSIS dhe Workflow nĂ« Informatica.

    Përveç dag-ve, gjithashtu mund të ketë sub-dag, por ndoshta nuk do të arrijmë atje.

  • DAG Run — njĂ« dag i inicializuar, tĂ« cilit i Ă«shtĂ« dhĂ«nĂ« e tij execution_date. DAGRAT e njĂ« dag-u mund tĂ« punojnĂ« paralelisht (nĂ«se, sigurisht, keni bĂ«rĂ« detyrat tuaja idempotente).
  • Operator — kĂ«to janĂ« copa kode qĂ« janĂ« pĂ«rgjegjĂ«se pĂ«r ekzekutimin e njĂ« veprimi tĂ« caktuar. Ka tre lloje operatorĂ«sh:
    • veprim, si pĂ«r shembull operatori ynĂ« i preferuar PythonOperator, i cili mund tĂ« ekzekutojĂ« çdo (tĂ« vlefshĂ«m) kod Python;
    • transfer, tĂ« cilĂ«t transportojnĂ« tĂ« dhĂ«na nga njĂ« vend nĂ« njĂ« tjetĂ«r, pĂ«r shembull, MsSqlToHiveTransfer;
    • sensor do tĂ« lejojĂ« tĂ« reagoni ose tĂ« ngadalĂ«soni ekzekutimin e dag-ut deri nĂ« ndodhjen e njĂ« eventi tĂ« caktuar. HttpSensor mund tĂ« tĂ«rheqĂ« endpoint-in e caktuar dhe kur tĂ« presĂ« pĂ«rgjigjen e duhur, tĂ« fillojĂ« transferimin GoogleCloudStorageToS3Operator. NjĂ« mendje kurioze do tĂ« pyesĂ«: "pĂ«rse? Sepse mund tĂ« bĂ«jmĂ« ripĂ«rsĂ«ritje direkt nĂ« operator!" Dhe pastaj, qĂ« tĂ« mos bllokojmĂ« rezervuarin e detyrave me operatorĂ« qĂ« janĂ« pezull. SensorĂ«t aktivizohen, kontrollojnĂ« dhe vdesin deri nĂ« pĂ«rpjekjen e ardhshme.
  • DetyrĂ« — operatorĂ«t e shpallur, pa marrĂ« parasysh llojin, dhe tĂ« lidhur me dagun rriten nĂ« nivelin e detyrave.
  • Instanca e detyrĂ«s — kur gjenerali-planifikues vendos se detyrat janĂ« gati pĂ«r t'u dĂ«rguar nĂ« luftĂ« te punĂ«torĂ«t-ekzekutorĂ« (nĂ« vend, nĂ«se po pĂ«rdorim LocalExecutor ose nĂ« njĂ« nod tĂ« largĂ«t nĂ« rastin e CeleryExecutor), ai u cakton atyre njĂ« kontekst (dmth. njĂ« grup variablesh - parametrash ekzekutimi), zhvillon modelet e komandave ose kĂ«rkesave dhe i vendos ato nĂ« rezervuar.

Drejtojmë detyrat

Së pari do të shënojmë skemën e përgjithshme të dagut tonë dhe pastaj do të thellojmë gjithnjë e më shumë në detaje, sepse përdorim disa zgjidhje jo triviale.

Pra, në formën më të thjeshtë, një dag i tillë do të dukej kështu:

nga datetime import timedelta, datetime

nga airflow import DAG
nga airflow.operators.python_operator import PythonOperator

nga 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)

për conn_id, schema në sql_server_ds:
    PythonOperator(
        task_id=schema,
        python_callable=workflow,
        provide_context=True,
        dag=dag)

Le të kuptojmë:

  • SĂ« pari importojmĂ« libraritĂ« e nevojshme dhe disa gjĂ«ra tĂ« tjera;
  • sql_server_ds — kjo Ă«shtĂ« List[namedtuple[str, str]] me emrat e lidhjeve nga Airflow Connections dhe bazat e tĂ« dhĂ«nave nga tĂ« cilat do tĂ« marrim tabelĂ«n tonĂ«;
  • dag — deklarata e DAG-ut tonĂ«, e cila duhet tĂ« jetĂ« nĂ« globals(), pĂ«rndryshe Airflow nuk do ta gjente. DAG-ut i duhet gjithashtu tĂ« themi:
    • se quhet orders — ky emĂ«r pastaj do tĂ« shfaqet nĂ« ndĂ«rfaqen e internetit,
    • se do tĂ« fillojĂ« punĂ«n nĂ« mesnatĂ« mĂ« 8 korrik,
    • dhe duhet tĂ« ekzekutohet, afĂ«rsisht çdo 6 orĂ« (pĂ«r djemtĂ« e mrekullueshĂ«m kĂ«tu nĂ« vend tĂ« timedelta() lejohet cron-linjĂ« 0 0 0/6 ? * * *, pĂ«r ata mĂ« pak tĂ« mrekullueshĂ«m — njĂ« shprehje si @daily);
  • workflow() do tĂ« bĂ«jĂ« punĂ«n kryesore, por jo tani. Tani ne thjesht do tĂ« hedhim kontekstin tonĂ« nĂ« log.
  • Dhe tani magjia e thjeshtĂ« e krijimit tĂ« detyrave:
    • ecim pĂ«rmes burimeve tona;
    • inicjalizojmĂ« PythonOperator, i cili do tĂ« realizojĂ« shoshin tonĂ« workflow(). Mos u harroni tĂ« tregoni njĂ« emĂ«r unik (brenda DAG-ut) pĂ«r detyrĂ«n dhe tĂ« lidheni me vetĂ« DAG-un. Flag siguroni_kontekst nga ana e tij do tĂ« hedhĂ« nĂ« funksion argumente shtesĂ«, tĂ« cilat ne do t'i mbledhim me kujdes pĂ«rmes **konteksti.

Derisa në këtë pikë, çfarë kemi marrë:

  • njĂ« DAG tĂ« ri nĂ« ndĂ«rfaqen web,
  • njĂ«qind e pesĂ«dhjetĂ« detyra qĂ« do tĂ« ekzekutohen paralelisht (nĂ«se lejojnĂ« konfigurimet e Airflow, Celery dhe fuqia e serverĂ«ve).

Të paktën pothuajse e morëm.

Apache Airflow: e bëjmë ETL më të lehtë
Kush do të instalohet varësitë?

Për ta thjeshtuar këtë, unë e kam ngjeshur në docker-compose.yml trajtim requirements.txt në të gjitha nodet.

Tani po fillon:

Apache Airflow: e bëjmë ETL më të lehtë

Kube gri — instanca tĂ« detyrave, tĂ« pĂ«rpunuara nga planifikuesi.

Disa minuta pritje, detyrat i kapin punëtorët:

Apache Airflow: e bëjmë ETL më të lehtë

TĂ« gjelbra, e qartĂ«, — qĂ« kanĂ« punuar me sukses. TĂ« kuqe — jo shumĂ« me sukses.

MeqenĂ«se, nĂ« prodhim tonĂ« nuk ka ndonjĂ« dosje ./dags, e cila sinkronizohet midis makinave — tĂ« gjitha DAG-tĂ« ndodhen nĂ« git nĂ« Gitlab tonĂ«, dhe Gitlab CI shpĂ«rndani pĂ«rditĂ«simet nĂ« makina kur behet bashkimi nĂ« master.

Pak për Flower

NdĂ«rsa punĂ«torĂ«t po punojnĂ« me detyrat tona boshe, kujtojmĂ« njĂ« mjet tjetĂ«r, i cili mund tĂ« na tregojĂ« diçka — Flower.

Faqja e parë me informacione përmbledhëse për nodet-punëtorë:

Apache Airflow: e bëjmë ETL më të lehtë

Faqja më e ngarkuar me detyra që kanë filluar punën:

Apache Airflow: e bëjmë ETL më të lehtë

Faqja më e mërzitshme me gjendjen e brokerit tonë:

Apache Airflow: e bëjmë ETL më të lehtë

Faqja mĂ« e ndritshme — me grafikat e gjendjes sĂ« detyrave dhe kohĂ«n e pĂ«rfundimit tĂ« tyre:

Apache Airflow: e bëjmë ETL më të lehtë

Po ngarkojmë të ngarkuarin

Pra, të gjithë detyrat kanë përfunduar, mund të çojmë të plagosurit.

Apache Airflow: e bëjmë ETL më të lehtë

Dhe numri i tĂ« plagosurve ishte i konsiderueshĂ«m — pĂ«r arsye tĂ« ndryshme. NĂ«se pĂ«rdoret siç duhet Airflow, kĂ«ta katrorĂ« tregojnĂ« se tĂ« dhĂ«nat me siguri nuk arritĂ«n.

Duhet të shohim logun dhe të rilançojmë instancat e detyrave që kanë dështuar.

Duke klikuar në çdo katror, do të shohim veprimet që na janë në dispozicion:

Apache Airflow: e bëjmë ETL më të lehtë

Mund të marrim dhe të bëjmë Clear për atë që ka dështuar. Pra, harrojmë se ka ndodhur diçka, dhe e njëjta instancë detyre do të shkojë te planifikuesi.

Apache Airflow: e bëjmë ETL më të lehtë

Natyrisht, nuk Ă«shtĂ« shumĂ« humane tĂ« bĂ«sh kĂ«shtu me tĂ« gjitha katrorĂ«t e kuq — nuk Ă«shtĂ« ky pritja jonĂ« nga Airflow. Sigurisht, kemi armĂ«n e masave shkatĂ«rruese: Browse/Task Instances

Apache Airflow: e bëjmë ETL më të lehtë

Do të zgjedhim gjithçka njëherësh dhe do të klikojmë opsionin e duhur për të kaluar në zero:

Apache Airflow: e bëjmë ETL më të lehtë

Pas pastrimit, taksitë tona duken kështu (ato tashmë po presin me padurim që planifikuesi t'i planifikojë):

Apache Airflow: e bëjmë ETL më të lehtë

Lidhje, likuj dhe variabla të tjerë

Koha më e mirë për të parë DAG-un tjetër, 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='Raportet DWH janë përditësuar',
    html_content=dedent("""Të nderuar, raportet janë përditësuar"""),
    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, zgjohem, ne {{ dag.dag_id }} e hodhëm
        """),
    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]

Çdokush ka bĂ«rĂ« ndonjĂ«herĂ« njĂ« raport azhurnimi, apo jo? Ja pĂ«rsĂ«ri: ka njĂ« listĂ« burimesh nga ku tĂ« marrim tĂ« dhĂ«nat; ka njĂ« listĂ« ku t'i vendosim; mos harro tĂ« na lajmerosh kur ndodhi ndonjĂ« gjĂ« ose kur diçka dĂ«shtoi (mirĂ«, kjo nuk Ă«shtĂ« pĂ«r ne, jo).

Le të kalojmë përsëri në skedarin dhe të shikojmë gjërat e reja enigmatike:

  • from commons.operators import TelegramBotSendMessage — asgjĂ« nuk na pengon tĂ« krijojmĂ« operatorĂ«t tanĂ«, siç e bĂ«mĂ« ne duke krijuar njĂ« mbĂ«shtjellĂ«s tĂ« vogĂ«l pĂ«r dĂ«rgimin e mesazheve nĂ« Razblokirovannyy. (PĂ«r kĂ«tĂ« operator do tĂ« flasim mĂ« poshtĂ«);
  • default_args={} — dag-u mund tĂ« ndajĂ« argumentet e njĂ«jta me tĂ« gjithĂ« operatorĂ«t e tij;
  • to='{{ var.value.all_the_kings_men }}' — fusha to nuk do tĂ« jetĂ« e koduar, por do tĂ« formohet dinamikisht me anĂ« tĂ« Jinja dhe njĂ« variabli me listĂ«n e email-eve qĂ« e kam vendosur me kujdes nĂ« Admin/Variables;
  • trigger_rule=TriggerRule.ALL_SUCCESS — kushti pĂ«r aktivizimin e operatorit. NĂ« rastin tonĂ«, letra do tĂ« shkojĂ« te shefat vetĂ«m nĂ«se tĂ« gjitha varĂ«sitĂ« pĂ«rfunduan suksesshĂ«m;
  • tg_bot_conn_id='tg_main' — argumentet conn_id pranojnĂ« identifikues tĂ« lidhjeve qĂ« krijojmĂ« nĂ« Admin/Connections;
  • trigger_rule=TriggerRule.ONE_FAILED — mesazhet nĂ« Telegram do tĂ« dĂ«rgohen vetĂ«m nĂ« rast se ka detyra tĂ« dĂ«shtuar;
  • task_concurrency=1 — ndalojmĂ« ekzekutimin e shumĂ« instanceve tĂ« njĂ«jta tĂ« njĂ« detyre nĂ« tĂ« njĂ«jtĂ«n kohĂ«. NĂ« tĂ« kundĂ«rt, do tĂ« kemi njĂ« ekzekutim tĂ« shumĂ« instancave VerticaOperator (tĂ« shikojnĂ« njĂ« tabelĂ«);
  • report_update >> [email, tg] — tĂ« gjitha VerticaOperator do tĂ« pĂ«rfundojnĂ« duke dĂ«rguar e-mail dhe mesazhe, kĂ«shtu:
    Apache Airflow: e bëjmë ETL më të lehtë

    Por, pasi operatorët e njoftimeve kanë kushte të ndryshme për t'u ekzekutuar, do të punojë vetëm njëri. Në Pamjen e Dendrës duket pak më pak e qartë:
    Apache Airflow: e bëjmë ETL më të lehtë

Do tĂ« them disa fjalĂ« pĂ«r makrotĂ« dhe miqtĂ« e tyre — variablat.

Makrotë janë placeholder-e Jinja, të cilat mund të vendosin informacion të ndryshëm të dobishëm në argumentet e operatorëve. Për shembull, kështu:

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

{{ ds }} do të zhvillohet në përmbajtjen e variablit të kontekstit execution_date në formatin YYYY-MM-DD: 2020-07-14. E bukura është se variablat e kontekstit lidhen fort me një instancë të caktuar të detyrës (katrorin në Pamjen e Dendrës), dhe kur bëhet një rinisje, placeholder-et do të zbulohen në të njëjtat vlera.

Vlerat e ndara mund të shikohen me butonin Rendered në çdo instancë të detyrës. Kështu është për detyrën me dërgimin e e-mailit:

Apache Airflow: e bëjmë ETL më të lehtë

Dhe kështu është për detyrën me dërgimin e mesazhit:

Apache Airflow: e bëjmë ETL më të lehtë

Lista e plotë e makrosve të integruar për versionin më të fundit të disponueshëm është këtu: Referenca për Makros

Për më tepër, me ndihmën e plugin-eve, ne mund të shpallim makros tona të personalizuara, por kjo është një histori krejt tjetër.

Përveç gjërave të paracaktuar, ne mund të përdorim vlerat e variablave tanë (më sipër në kod unë e kam përdorur këtë). Le të krijojmë në Admin/Variables një çift gjërash:

Apache Airflow: e bëjmë ETL më të lehtë

Gati pĂ«r t’u pĂ«rdorur:

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

Në vlerë mund të jetë skalar, ose mund të jetë edhe JSON. Në rastin e JSON-it:

bot_config

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

thjesht përdorim rrugën për çelësin e kërkuar: {{ var.json.bot_config.bot.token }}.

Do të them vetëm një fjalë dhe do të tregoj një skrin që ka të bëjë me lidhet. Këtu gjithçka është elementare: në faqe Admin/Connections krijojmë lidhjen, vendosim atje emrat e përdoruesve/fjalëkalimet dhe parametra më specifik. Ja kështu:

Apache Airflow: e bëjmë ETL më të lehtë

FjalĂ«kalimat mund tĂ« shkruhen (mĂ« me kujdes se nĂ« variantin e paracaktuar), ose mund tĂ« mos caktohet lloji i lidhjes (ashtu siç bĂ«ra pĂ«r tg_main) — çështja Ă«shtĂ« se lista e llojeve Ă«shtĂ« e koduar nĂ« modelet Airflow dhe zgjerimi nuk lejon ndryshime nĂ« burim (nĂ«se ndonjĂ«herĂ« nuk kam gjetur diçka — ju lutem mĂ« korrigjoni), por nuk ka asgjĂ« qĂ« na ndalon tĂ« marrim kredencialet thjesht me emrin.

Gjithashtu, mund të bëni disa lidhje me një emër të vetëm: në këtë rast, metoda BaseHook.get_connection(), e cila na merr lidhjet sipas emrit, do të kthejë një rastësor nga disa të ngjashëm (do të ishte më logjike të bëhej Round Robin, por do ta lëmë këtë në ndërgjegjen e zhvilluesve të Airflow).

Variables dhe Connections, pa dyshim, janĂ« mjete tĂ« shkĂ«lqyera, por Ă«shtĂ« e rĂ«ndĂ«sishme tĂ« mos humbasĂ«sh balancimin: cilat pjesĂ« tĂ« rrjedhave tuaja i mbani nĂ« kodin tuaj dhe cilat — i jepni pĂ«r ruajtje nĂ« Airflow. Nga njĂ«ra anĂ«, ndĂ«rrimi i shpejtĂ« i vlerĂ«s, pĂ«r shembull, box-i i dĂ«rgesĂ«s, mund tĂ« jetĂ« i lehtĂ« pĂ«rmes UI. Nga ana tjetĂ«r, kjo Ă«shtĂ« nĂ« fund tĂ« fundit njĂ« kthim nĂ« klikimin me maus, nga i cili ne (unĂ«) do doja tĂ« shpenzoja.

Puna me lidhjet — Ă«shtĂ« njĂ« nga detyrat hook-esh. NĂ« pĂ«rgjithĂ«si, hook-at e Airflow janĂ« pika lidhĂ«se me shĂ«rbime dhe biblioteka tĂ« jashtme. P.sh., JiraHook do tĂ« hapĂ« pĂ«r ne njĂ« klient pĂ«r tĂ« bashkĂ«punuar me Jira (mund tĂ« lĂ«vizim detyra kĂ«tu-kĂ«tu), dhe me ndihmĂ«n e SambaHook mund tĂ« dĂ«rgoni njĂ« skedar lokal nĂ« smb-pikĂ«.

Analizojmë operatorin e personalizuar

Dhe ne u afruam ngushtësisht për të parë se si është bërë TelegramBotSendMessage

Kodi commons/operators.py me operatorin vetë:

from typing import Union

from airflow.operators import BaseOperator

from commons.hooks import TelegramBotHook, TelegramBot

class TelegramBotSendMessage(BaseOperator):
    """Dërgo mesazh në chat_id duke përdorur TelegramBotHook

    Shembull:
        >>> 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 }} dështoi :(',
        ...     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'Dërgo "{self.message}" në bisedën {self.chat_id}')
        self.client.send_message(chat_id=self.chat_id,
                                 message=self.message)

Këtu, ashtu si gjithçka tjetër në Airflow, është shumë e thjeshtë:

  • TrashĂ«guam nga BaseOperator, e cila implementon shumĂ« gjĂ«ra specifike pĂ«r Airflow (shikoni kur keni kohĂ«)
  • Deklaruam fushat template_fields, nĂ« tĂ« cilat Jinja do tĂ« kĂ«rkojĂ« makro pĂ«r pĂ«rpunim.
  • Organizuam argumentet e duhura pĂ«r __init__(), vendosĂ«m default-et atje ku duhet.
  • Nuk harrojmĂ« as inicializimin e prindit.
  • HapĂ«m hook-un pĂ«rkatĂ«s TelegramBotHook, morĂ«m prej tij njĂ« objekt-klient.
  • E overrajduam (pĂ«rkufizuam pĂ«rsĂ«ri) metodĂ«n BaseOperator.execute(), qĂ« Airflow do ta thĂ«rrasĂ« kur Ă«shtĂ« koha pĂ«r tĂ« filluar operatorin — aty do tĂ« realizojmĂ« veprimin kryesor, pa harruar tĂ« regjistrojmĂ«. (Regjistrohemi, pĂ«r shembull, direkt nĂ« stdout dhe stderr — Airflow do ta kapĂ« gjithçka, do ta mbĂ«shtjellĂ« bukur, do ta vendosĂ« ku duhet.)

Të shikojmë çfarë kemi në commons/hooks.py. Pjesa e parë e skedarit, me vetë hook-un:

from typing import Union

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

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

    Shënim: shtoni një lidhje me Conn Type bosh dhe mos harroni
    të plotësoni 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

S'di ashtë çfarë mund të shpjegohet këtu, thjesht do të theksoj disa pika të rëndësishme:

  • Ne trashĂ«gojmĂ«, mendojmĂ« pĂ«r argumentet — nĂ« shumicĂ«n e rasteve do tĂ« jetĂ« vetĂ«m njĂ«: conn_id;
  • TejkalojmĂ« metodat standarde: unĂ« e kufizova get_conn(), ku marrĂ« parametrat e lidhjes sipas emrit dhe thjesht nxjerr seksionin extra (ky Ă«shtĂ« njĂ« fushĂ« pĂ«r JSON), nĂ« tĂ« cilĂ«n unĂ« (sipas udhĂ«zimit tim!) vendosa tokenin e bosit Telegram: {"bot_token": "YOuRAwEsomeBOtToKen"}.
  • Krijoj njĂ« instancĂ« tĂ« TelegramBot, duke i dhĂ«nĂ« gjithashtu tokenin e saktĂ«.

Këtu është. Të marrë klientin nga huka mund të bëhet me TelegramBotHook().clent ose TelegramBotHook().get_conn().

Dhe pjesa e dytë e skedarit, ku bëj një mikro mbështjellje për Telegram REST API, për të mos sjellur të njëjtin python-telegram-bot për një metodë të vetme 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))

QĂ«llimi i duhur Ă«shtĂ« tĂ« bashkosh tĂ« gjitha kĂ«to: TelegramBotSendMessage, TelegramBotHook, TelegramBot — nĂ« plugin, ta vendosĂ«sh nĂ« njĂ« depo tĂ« hapur dhe ta ofrosh si Open Source.

Përderisa po e shqyrtonim gjithçka, përditësimet tona të raporteve arritën që të dërgojnë një mesazh gabimi në kanalin tim. Po shkoj të kontrolloj se çfarë është sërish gabim...

Apache Airflow: e bëjmë ETL më të lehtë
Diçka është prishur në dagun tonë! Po ashtu kjo është ajo që prisnim? Pikërisht!

A do të derdhësh?

Ndiheni se kam harruar diçka? Duket se premtova të transferoj të dhënat nga SQL Server në Vertica dhe këtu u largova nga tema, sa bezdisës!

Krimi ishte i qëllimtë, thjesht kisha detyrim të shpjegoja disa terminologji për ju. Tani mund të vazhdojmë.

Plani ynë ishte si më poshtë:

  1. Të bëjmë dag
  2. Të gjenerojmë detyra
  3. Të shohim si duket gjithçka
  4. Të caktojmë numrat e seancave për mbushjet
  5. Të marrim të dhënat nga SQL Server
  6. Të vendosim të dhënat në Vertica
  7. Të grumbullojmë statistikat

Pra, për të gjithë këtë, kam bërë një shtesë të vogël në 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

Atje shohim:

  • Vertica si host dwh me cilĂ«simet mĂ« tĂ« zakonshme,
  • tri instanca SQL Server,
  • plotsojmĂ« bazat me disa tĂ« dhĂ«na tĂ« fundit (nuk duhet tĂ« shikoni nĂ« mssql_init.py!)

E kemi nisur gjithçka me një komandë pak më të komplikuar se më parë:

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

ÇfarĂ« ka gjeneruar rastĂ«sia jonĂ«, mund tĂ« shikohet duke pĂ«rdorur pikĂ«n Profilizimi i tĂ« DhĂ«nave/KĂ«rkesa Ad Hoc:

Apache Airflow: e bëjmë ETL më të lehtë
Kryesore, mos e tregoni këtë analistëve

Të ndalemi në detaje seancat ETL nuk do ta bëj, aty është gjithçka e thjeshtë: krijojmë një bazë, në të një tabelë, e mbështjellim gjithçka me menaxherin e kontekstit, dhe tani veprojmë kështu:

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

    # Load worflow
    ...

    session.successful = True
    session.loaded_rows = 15

session.py

nga sys import stderr

klasa Session:
    """Seanca e punës ETL

    Example:
        me Session(task_name) si sesion:
            print(session.id)
            session.successful = True
            session.loaded_rows = 15
            session.comment = 'Mirë
    """

    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 = """
            KRIJO TABELË NËSE NUK EKZISTON sessions (
                id          SERIAL       NUK NËN NORMËS
                task_name   VARCHAR(200) NUK NËN NORMËS,

                started     TIMESTAMPTZ  NUK NËN NORMËS DEFAULT current_timestamp,
                finished    TIMESTAMPTZ           DEFAULT current_timestamp,
                successful  BOOL,

                loaded_rows INT,
                comment     VARCHAR(500)
            );
            """
        self._execute(query)

    def open(self):
        query = """
            SHTO NË sessions (task_name, finished)
            VLERAT (%s, NULL)
            KTHIM id;
            """
        self._id = self._execute(query, self.task_name)
        print(self, 'hapur')
        return self

    def close(self):
        if not self._id:
            raise SessionClosedError('Seanca nuk është e hapur')
        query = """
            PËRMIKSO sessions
            SET
                finished    = DEFAULT,
                successful  = %s,
                loaded_rows = %s,
                comment     = %s
            KU
                id = %s
            KTHIM id;
            """
        self._execute(query, self.successful, self.loaded_rows,
                      self.comment, self.id)
        print(self, 'mbyllur',
              ', e suksesshme: ', self.successful,
              ', Ngarkuar: ', self.loaded_rows,
              ', komenti:', self.comment)

klasa SessionError(Exception):
    kalon

klasa SessionClosedError(SessionError):
    kalon

Ka ardhur koha të marrim të dhënat tona nga pesëmbëdhjetë tabela tona. Do ta bëjmë këtë me disa rreshta të thjeshtë:

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. Me anë të hook-ut do marrim nga Airflow pymssql-konnekt
  2. NĂ« kĂ«rkesĂ« do vendosim njĂ« kufizim nĂ« formĂ«n e datĂ«s — kĂ«tĂ« do ta vendosĂ« shabloni.
  3. TĂ«rheqim kĂ«rkesĂ«n tonĂ« pandas, e cila do tĂ« sjellĂ« pĂ«r ne DataFrame — do na nevojitet mĂ« vonĂ«.

Unë përdor zëvendësimin {dt} në vend të parametrave të kërkesës %s jo pse unë jam një Buratino i keq, por sepse pandas nuk mund të përballojë me pymssql dhe i ofron këtij të fundit params: List, ndonëse ai shumë dëshiron tuple.
Gjithashtu, vini re se zhvilluesi pymssql vendosi të mos e mbështesë më, dhe është koha të kalojmë në pyodbc.

Le të shohim çfarë i shtoi Airflow argumenteve tona:

Apache Airflow: e bëjmë ETL më të lehtë

Nëse nuk ka të dhëna, nuk ka kuptim të vazhdojmë. Por gjithashtu është çuditshëm të quhet ngarkimi i suksesshëm. Po ashtu, nuk është një gabim. Ah, çfarë të bëjmë?! Ja çfarë:

if df.empty:
    raise AirflowSkipException('Nuk ka rreshta për të ngarkuar')

AirflowSkipException do të thotë Airflow se nuk ka gabime, por ne do ta kalojmë detyrën. Në ndërfaqe do të ketë një katror që nuk është gjelber dhe as i kuq, por me ngjyrë rozë.

Do t'i hedhim të dhënat tona për disa kolona:

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

Specifikisht:

  • DB, nga e cila kemi marrĂ« porositĂ«,
  • Identifikuesi i sesionit tonĂ« tĂ« ngarkimit (do tĂ« jetĂ« i ndryshĂ«m pĂ«r secilĂ«n detyrĂ«),
  • Hash nga burimi dhe identifikuesi i porosisĂ« — qĂ« nĂ« bazĂ«n pĂ«rfundimtare (ku gjithçka do tĂ« bashkohet nĂ« njĂ« tabelĂ«) tĂ« kemi njĂ« identifikues unik pĂ«r porosinĂ«.

Ka mbetur hapi i parafundit: të ngarkohet gjithçka në Vertica. Dhe, siç është e çuditshme, një nga mënyrat më efektive për ta bërë këtë është përmes 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. Ne bëjmë një marrëveshje speciale StringIO.
  2. pandas do të përfshijë në të DataFrame në formën e rreshtat CSV.Do të hapim një lidhje me Vertica-n tonë të dashur.
  3. Dhe tani me ndihmën e
  4. copy() do t'i dërgojmë të dhënat tona direkt në Vertica! do t'i dërgojmë të dhënat tona direkt në Vertica!

Do të marrim nga driver-i sa rreshta janë ngarkuar, dhe do t'i themi menaxherit të sesionit që gjithçka është OK:

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

Këtu është gjithçka.

Në prodhim krijojmë tabelën e synuar manualisht. Këtu le të lejoj vetes një automatizëm të vogël:

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)

Unë me ndihmën e VerticaOperator() krijoj skemën e DB-së dhe tabelën (nëse ato nuk ekzistojnë, sigurisht). E rëndësishme është të vendosësh drejt varësitë:

për conn_id, schema në 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

Përmbledhje

— Ja, — tha minusi, — a nuk Ă«shtĂ« e vĂ«rtetĂ« se tani
Ti e ke kuptuar që ndodhem unë si kafsha më e tmerrshme në pyll?

Julia Donaldson, «Gruffalo»

Mendoj se po tĂ« organizonim njĂ« garĂ« mes kolegĂ«ve tĂ« mi: kush e pĂ«rgatit dhe e nis mĂ« shpejt ETL-procesin nga fillimi: ata me SSIS dhe miniku, dhe unĂ« me Airflow
 Pastaj do tĂ« krahasonim edhe lehtĂ«sinĂ« e mirĂ«mbajtjes
 Ehm, mendoj se do tĂ« pajtohesh se do t'i kaloj nĂ« çdo front!

NĂ«se flasim pak mĂ« seriozisht, atĂ«herĂ« Apache Airflow — pĂ«rmes pĂ«rshkrimit tĂ« proceseve nĂ« formĂ«n e kodit programues — e bĂ«ri punĂ«n time shumĂ« mĂ« mĂ« tĂ« lehtĂ« dhe mĂ« tĂ« kĂ«ndshme.

Po ashtu, shkallĂ«zueshmĂ«ria e tij e pakufizuar: si nĂ« aspektin e plugins, ashtu edhe prirja pĂ«r tĂ« shkallĂ«zuar — ju jep mundĂ«sinĂ« tĂ« pĂ«rdorni Airflow nĂ« pothuajse çdo fushĂ«: qoftĂ« nĂ« ciklin e plotĂ« tĂ« mbledhjes, pĂ«rgatitjes dhe pĂ«rpunimit tĂ« tĂ« dhĂ«nave, qoftĂ« nĂ« raketat qĂ« lĂ«shohen (nĂ« Mars, sigurisht).

Pjesa përfundimtare, informative dhe reference

Gabimet që i kemi mbledhur për ju

  • start_date. Po, kjo Ă«shtĂ« njĂ« meme lokale. PĂ«rmes argumentit kryesor tĂ« DAG-ut start_date kalojnĂ« tĂ« gjitha. Shkurtimisht, nĂ«se caktojmĂ« nĂ« start_date datĂ«n aktuale, dhe nĂ« schedule_interval — njĂ« ditĂ«, atĂ«herĂ« DAG-u do tĂ« nisĂ« nesĂ«r, mĂ« herĂ«t se.
    start_date = datetime(2020, 7, 7, 0, 1, 2)

    Dhe më shumë asnjë problem.

    Me të lidhet dhe një tjetër gabim në ekzekutim: Task is missing the start_date parameter, që zakonisht tregon se e keni harruar të lidhni me operatorin e DAG-ut.

  • TĂ« gjitha nĂ« njĂ« makinĂ«. Po, si bazat (e vetĂ« Airflow dhe mbulesĂ«s sonĂ«), ashtu edhe serveri web, planifikuesi dhe punĂ«torĂ«t. Dhe madje punonte. Por me kalimin e kohĂ«s numri i detyrave nĂ« shĂ«rbime u rrit, dhe kur PostgreSQL filloi tĂ« jepte pĂ«rgjigje nĂ«pĂ«rmjet indekseve pĂ«r 20 ms nĂ« vend tĂ« 5 ms, ne e morĂ«m dhe e trezoruam.
  • LocalExecutor. Po, ne kemi qenĂ« ende nĂ« tĂ«, dhe tashmĂ« kemi arritur nĂ« skajin e greminĂ«s. LocalExecutor ka qenĂ« i mjaftueshĂ«m pĂ«r ne deri tani, por tani ka ardhur koha pĂ«r tĂ« zgjeruar me tĂ« paktĂ«n njĂ« punĂ«tor, dhe do tĂ« duhet tĂ« mundohemi tĂ« kalojmĂ« nĂ« CeleryExecutor. Dhe duke marrĂ« parasysh se mund tĂ« punosh me tĂ« edhe nĂ« njĂ« makinĂ«, asgjĂ« nuk na ndalon nga pĂ«rdorimi i Celery edhe nĂ« serverin qĂ« "natyrisht, kurrĂ« nuk do tĂ« shkojĂ« nĂ« prodhim, tĂ« premtoj!"
  • PĂ«rdorim i papĂ«rfshirĂ« mjetet e ndĂ«rtuara:
    • Connections pĂ«r ruajtjen e tĂ« dhĂ«nave tĂ« identifikimit tĂ« shĂ«rbimeve,
    • SLA Misses pĂ«r tĂ« reaguar ndaj detyrave qĂ« nuk pĂ«rfunduan nĂ« kohĂ«,
    • XCom pĂ«r ndarjen e metadatat (thashĂ« metatĂ« dhĂ«nat!) midis detyrave tĂ« dagut.
  • Abuzimi me postĂ«n. ÇfarĂ« tĂ« thuash kĂ«tu? Ishin krijuar njoftime pĂ«r tĂ« gjitha pĂ«rsĂ«ritjet e detyrave tĂ« rĂ«na. Tani nĂ« Gmail-in tim tĂ« punĂ«s kam >90k letra nga Airflow, dhe ndĂ«rfaqja web e postĂ«s refuzon tĂ« marrĂ« dhe tĂ« fshijĂ« mĂ« shumĂ« se 100 njĂ«si pĂ«r herĂ«.

Më shumë pengesa: Apache Airflow Pitfalls

Mjetet për automatizim më të madh

Për të punuar edhe më shumë me mendje dhe jo me duar, Airflow na përgatiti këtë:

  • REST API — ai ende ka status Experimental, qĂ« nuk e pengon tĂ« funksionojĂ«. Me kĂ«tĂ«, jo vetĂ«m qĂ« mund tĂ« merrni informacion mbi dagat dhe detyrat, por mund tĂ« ndaloni/nisni dagun, tĂ« krijoni DAG Run ose grup.
  • CLI — pĂ«rmes komandĂ«s nĂ« terminal janĂ« tĂ« ndAvailable shumĂ« mjete, tĂ« cilat jo vetĂ«m qĂ« janĂ« tĂ« pakĂ«ndshme pĂ«r t'u pĂ«rdorur pĂ«rmes WebUI, por edhe nuk ekzistojnĂ« fare. PĂ«r shembull:
    • backfill nevojiten pĂ«r tĂ« rinisur instancat e detyrave.
      Për shembull, vijnë analistët dhe thonë: «E keni, shoku, problem me të dhënat nga 1 deri më 13 janar! Rregulloni!». Dhe ti je si:
      airflow backfill -s '2020-01-01' -e '2020-01-13' orders
    • MirĂ«mbajtja e bazĂ«s: initdb, resetdb, upgradedb, checkdb.
    • run, i cili lejon tĂ« nisni njĂ« instancĂ« tĂ« detyrĂ«s dhe t'i injoroni tĂ« gjitha varĂ«sitĂ«. PĂ«r mĂ« tepĂ«r, mund ta nisni atĂ« pĂ«rmes LocalExecutor, madje edhe sikur tĂ« keni njĂ« klaster Celery.
    • PĂ«rgjithĂ«sisht, e njĂ«jta gjĂ« bĂ«nĂ« test, vetĂ«m se nuk shkruan asgjĂ« nĂ« bazĂ«.
    • connections lejon krijimin masiv tĂ« lidhjeve nga shell.
  • Python API — njĂ« mĂ«nyrĂ« mjaft hardcore pĂ«r tĂ« bashkĂ«vepruar, e cila Ă«shtĂ« e destinuar pĂ«r plugina, dhe jo pĂ«r tĂ« grumbulluar me duar. Por kush na ndalon tĂ« shkojmĂ« nĂ« /home/airflow/dags, tĂ« nxjerrim ipython dhe tĂ« fillojmĂ« tĂ« bĂ«jmĂ« çfarĂ« deshi? Mund, pĂ«r shembull, tĂ« eksportoni tĂ« gjitha lidhjet me kĂ«tĂ« kod:
    nga airflow import settings
    nga airflow.models import Connection
    
    fushat = 'conn_id conn_type host port schema login password extra'.split()
    
    session = settings.Session()
    per conn në session.query(Connection).order_by(Connection.conn_id):
      d = {field: getattr(conn, field) për fushat}
      print(conn.conn_id, '=', d)
  • Lidhja me bazĂ«n e tĂ« dhĂ«nave tĂ« metadatas Airflow. Nuk e rekomandoj tĂ« shkruani nĂ« tĂ«, por merrni statet e detyrave pĂ«r metrika tĂ« ndryshme specifike shumĂ« mĂ« shpejt dhe mĂ« lehtĂ« se pĂ«rmes ndonjĂ« API.

    TĂ« themi, jo tĂ« gjitha detyrat tona janĂ« idempotente, dhe ndonjĂ«herĂ« ato mund tĂ« dĂ«shtojnĂ«, dhe kjo Ă«shtĂ« e pranueshme. Por disa dĂ«shtime — kjo Ă«shtĂ« tashmĂ« e dyshimtĂ«, dhe duhet tĂ« verifikohet.

    Kujdes, SQL!

    ME ekzekutimet e fundit SI (
    SELECT
        task_id,
        dag_id,
        execution_date,
        state,
            row_number()
            PËR (
                PARTITION BY task_id, dag_id
                ORDER BY execution_date DESC) SI rn
    FROM public.task_instance
    WHERE
        execution_date > now() - INTERVAL '2' DITË
    ),
    failed AS (
        SELECT
            task_id,
            dag_id,
            execution_date,
            state,
            RAST CASE kur rn = row_number() PËR (
                PARTITION BY task_id, dag_id
                ORDER BY execution_date DESC)
                     THEN TRUE END SI last_fail_seq
        FROM last_executions
        WHERE
            state IN ('dështuar', 'në pritje për rikthim')
    )
    SELECT
        task_id,
        dag_id,
        count(last_fail_seq)                       SI të pasuksesshëm,
        count(CASE KUR last_fail_seq
            DHE state = 'dĂ«shtuar' ATËHERË 1 END)       SI dĂ«shtuar,
        count(CASE KUR last_fail_seq
            DHE state = 'nĂ« pritje pĂ«r rikthim' ATËHERË 1 END) SI nĂ« pritje pĂ«r rikthim
    FROM failed
    GRUPI PËR
        task_id,
        dag_id
    KUSH
        count(last_fail_seq) > 0

Linke

Natyrisht, dhjetë lidhjet e para nga rezultatet e Google përmbajnë dosjen e Airflow nga shënimet e mia.

Dhe lidhjet e përdorura në artikull:

Burimi: habr.com

Bli njĂ« hosting tĂ« besueshĂ«m pĂ«r faqet me mbrojtje DDoS, VPS VDS serverĂ« đŸ”„ Bli njĂ« hosting tĂ« besueshĂ«m pĂ«r faqet me mbrojtje DDoS, VPS VDS serverĂ« | ProHoster