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

PĂ«rshĂ«ndetje, unĂ« jam Dmitriy Logvinenko — Inxhinier tĂ« DhĂ«nash nĂ« departamentin e analitikĂ«s sĂ« grupit tĂ« kompanive «Veze».

Do t'ju flas pĂ«r njĂ« mjet tĂ« shkĂ«lqyer pĂ«r zhvillimin e proceseve ETL — Apache Airflow. Por Airflow Ă«shtĂ« aq universale dhe shumĂ«anshĂ«m, sa duhet t'i kushtoni vĂ«mendje edhe nĂ«se nuk merresh me flukse tĂ« tĂ« dhĂ«nash, por keni nevojĂ« pĂ«r tĂ« nisur periodicisht ndonjĂ« proces dhe tĂ« monitoroni pĂ«rmbushjen e tij.

Dhe po, nuk do të flas vetëm, por do të tregoj gjithashtu: programi përmban shumë kod, screenshot-e dhe rekomandime.

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

Përmbajtja

Hyrje

Apache Airflow — Ă«shtĂ« si Django:

  • Ă«shtĂ« shkruar nĂ« Python,
  • ka njĂ« panel administrimi tĂ« shkĂ«lqyer,
  • Ă«shtĂ« pa kufi nĂ« zgjerim,

— vetĂ«m mĂ« mirĂ«, dhe Ă«shtĂ« bĂ«rĂ« krejtĂ«sisht pĂ«r qĂ«llime tĂ« tjera, siç Ă«shtĂ« shkruar deri nĂ« kata:

  • nĂ« ekzekutimin dhe monitorimin e detyrave nĂ« njĂ« numĂ«r tĂ« pakufizuar makinash (sa ju lejon Celery/Kubernetes dhe ndjenja juaj e pĂ«rgjegjĂ«sisĂ«)
  • me gjenerimin dinamik tĂ« punĂ«ve nga njĂ« kod Python shumĂ« tĂ« lehtĂ« pĂ«r t'u shkruar dhe kuptuar
  • dhe me mundĂ«sinĂ« pĂ«r tĂ« lidhur çdo bazĂ« tĂ« dhĂ«nash dhe API me ndihmĂ«n e komponenteve tĂ« gatshme, si dhe plugina tĂ« bĂ«ra nĂ« shtĂ«pi (tĂ« cilat janĂ« shumĂ« tĂ« lehta pĂ«r t'u bĂ«rĂ«).

Ne përdorim Apache Airflow kështu:

  • mbledhim tĂ« dhĂ«na nga burime tĂ« ndryshme (marrĂ«dhĂ«nie tĂ« shumta SQL Server dhe PostgreSQL, API tĂ« ndryshme me metrikat e aplikacioneve, madje edhe 1C) nĂ« DWH dhe ODS (pĂ«r ne Ă«shtĂ« Vertica dhe Clickhouse).
  • si njĂ« avansuar cronqĂ« nxit proceset e konsolidimit tĂ« tĂ« dhĂ«nave nĂ« ODS, si dhe monitoron shĂ«rbimin e tyre.

Derisa një kohë të shkurtër më parë nevojat tona u ngopën nga një server i vogël me 32 CPU dhe 50 GB RAM. Në Airflow, për këtë arsye, punon:

  • mĂ« shumĂ« se 200 DAG-Ă« (nĂ« fakt flukset e punĂ«s, nĂ« tĂ« cilat ne kemi mbushur detyrat),
  • nĂ« mesin e secilit mesatarisht 70 detyrash,
  • kjo gjĂ« pĂ«rfundohet (poashtu mesatarisht) njĂ« herĂ« nĂ« orĂ«.

Dhe për atë se si ne u zgjeruam, do të shkruaj më poshtë, por tani le t'ia përcaktojmë përmendësh-qartë detyrën që do të zgjidhim:

Ka there tre SQL Server përdorues, çdo një me 50 baza të dhënash - instanca të një projekti, për pasojë, struktura e tyre është identike (gati kudo, muahaha), dhe do të thotë se në çdo njëri ka një tabelë të quajtur Orders (me fat që një tabelë me këtë emër mund të vendoset në çdo biznes). Ne marrim të dhënat, duke shtuar fushat ndihmëse (serveri burim, baza burim, identifikues i detyrës ETL) dhe naivisht do t'i hedhin ato në, le të themi, Vertica.

Le të shkojmë!

Pjesa kryesore, praktike (dhe pak teorike)

Pse na nevojitet ajo (dhe juve)

Kur pemët ishin të mëdha dhe unë isha thjesht SQL-punëtor në një shitës me pakicë rus, ne po zhvillonim proceset ETL aka rrjedhat e të dhënave me dy mjetet që kishim në dispozicion:

  • Informatica Power Center — njĂ« sistem tepĂ«r i hollĂ«sishĂ«m, jashtĂ«zakonisht i fuqishĂ«m, me harduerin e tij, versionimin e tij. Kam pĂ«rdorur, mĂ« mirĂ« tĂ« them, 1% tĂ« mundĂ«sive tĂ« tij. Pse? SĂ« pari, ky ndĂ«rfaqe duket si nga fillimi i viteve 2000 dhe na jashtĂ«zakonisht ushtroi presion. SĂ« dyti, ky gjĂ« Ă«shtĂ« e projektuar pĂ«r procese shumĂ« tĂ« ndĂ«rlikuara, ri-pĂ«rdorimin intensiv tĂ« komponentĂ«ve dhe tĂ« tjera gjĂ«ra tĂ« rĂ«ndĂ«sishme pĂ«r ndĂ«rmarrjet. PĂ«r atĂ« qĂ« kushton, si krah Airbus A380/ vit, ne do tĂ« heshtim.

    Kujdes, një screenshot mund t'i dëmtojë të rinjtë nën 30 vjeç.

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

  • SQL Server Integration Server — ky shok e kemi pĂ«rdorur nĂ« rrjedhat tona brenda projektit. Ajo qĂ« Ă«shtĂ« e vĂ«rtetĂ«: ne tashmĂ« po pĂ«rdorim SQL Server dhe tĂ« mos pĂ«rdorim mjete ETL tĂ« tij do tĂ« ishte disi e paarsyeshme. Lidhur me tĂ« gjitha, dhe ndĂ«rfaqja Ă«shtĂ« e bukur, dhe raportet e performancĂ«s
 Por nuk e duam kĂ«tĂ« pĂ«r produktet e softuerit, oh jo, jo pĂ«r kĂ«tĂ«. Ne mund tĂ« versionojmĂ« tĂ« dtsx (i cili Ă«shtĂ« njĂ« XML me node qĂ« pĂ«rzihen gjatĂ« ruajtjes) ne mundemi, por çfarĂ« dobi? A tĂ« bĂ«jmĂ« njĂ« paketĂ« detyrash e cila do tĂ« transferonte qindra tabela nga njĂ« server nĂ« tjetrin? Ta zĂ«mĂ« se jo qindra, por 20 do t'ju bĂ«nte tĂ« humbni gishti tregues qĂ« klikoni mbi butonin e mausit. Por pamja e saj, padyshim Ă«shtĂ« mĂ« moderne:

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

Patjetër që po kërkonim zgjidhje. Madje gati arritëm deri tek një gjenerator i vetë-shkruar të paketimeve SSIS


... dhe pastaj më gjeti një punë e re. Aty më ndali Apache Airflow.

Kur zbulova se përshkrimi i proceseve ETL është një thjesht kod Python, nuk kisha parë ndonjëherë një gëzim të tillë. Ja, si u versionuan rrjedhat e të dhënave dhe u shkarkuan, dhe derdhja e tabelave me strukturë të përbashkët nga qindra baza të dhënash në një target bëri që të ishte një punë e thjeshtë e kodit Python në ekranet e 13'.

Krijojmë një klaster

Mos u bëjmë një kopësht të fëmijëve, dhe të flasim këtu për gjëra të qarta, si instalimi i Airflow, baza e të dhënave që keni zgjedhur, Celery dhe punë të tjera që përmenden në dokumenta.

Që të mund të fillojmë menjëherë me eksperimente, unë kam shkruar docker-compose.yml në të cilin:

  • Do tĂ« ngrisim me tĂ« vĂ«rtetĂ« Airflow: Planifikuesi, Serveri i Webit. Atje gjithashtu do tĂ« funksionojĂ« Flower pĂ«r monitorimin e detyrave tĂ« Celery (sepse tashmĂ« Ă«shtĂ« futur nĂ« apache/airflow:1.10.10-python3.7, dhe ne s'kemi asgjĂ« kundĂ«r);
  • PostgreSQL, nĂ« tĂ« cilin Airflow do tĂ« shkruajĂ« informacionin e tij operativ (tĂ« dhĂ«nat e planifikuesit, statistikat e ekzekutimit etj.), ndĂ«rsa Celery do tĂ« regjistrojĂ« detyrat e pĂ«rfunduara;
  • Redis, i cili do tĂ« veprojĂ« si broker pĂ«r detyrat e Celery;
  • PunĂ«tor Celery, i cili do tĂ« merret me ekzekutimin e detyrave direkt.
  • NĂ« dosjen . /dags ne do tĂ« vendosim skedarĂ«t tanĂ« me pĂ«rshkrimin e DAG-ve. Ato do tĂ« kapen nĂ« flakĂ«, kĂ«shtu qĂ« nuk Ă«shtĂ« e nevojshme tĂ« çaktivizoni tĂ«rĂ« sistemin pas çdo gjesti tĂ« vogĂ«l.

Disa kode në shembuj nuk janë të plota (për të mos e ngarkuar tekstin), dhe ndonjëherë ato modifikohen gjatë procesit. Shembujt e plotë të punës mund të shikohen në depo. 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

Vërejtje:

  • NĂ« ndĂ«rtimin e kompozitĂ«s, unĂ« u mbĂ«shteta shumĂ« nĂ« imazhin e njohur puckel/docker-airflow – sigurisht qĂ« duhet ta shikoni. Ndoshta nuk do t'ju nevojitet mĂ« asgjĂ« tjetĂ«r nĂ« jetĂ«.
  • TĂ« gjitha konfigurimet e Airflow janĂ« tĂ« disponueshme jo vetĂ«m pĂ«rmes airflow.cfg, por edhe pĂ«rmes variablave tĂ« mjedisit (falĂ« zhvilluesve), qĂ« unĂ« e shfrytĂ«zova armiqĂ«sisht.
  • Natyrisht, ai nuk Ă«shtĂ« pĂ«r prodhim: e kam bĂ«rĂ« qĂ«llimisht pa heartbeats nĂ« kontejnerĂ«t, dhe nuk kam shqetĂ«suar pĂ«r sigurinĂ«. Por minimumi, i pĂ«rshtatshĂ«m pĂ«r eksperimentet tona, e kam bĂ«rĂ«.
  • Vini re se:
    • Folderi me DAG duhet tĂ« jetĂ« i aksesueshĂ«m si pĂ«r planifikuesin ashtu edhe pĂ«r punĂ«torĂ«t.
    • E njĂ«jta gjĂ« vlen pĂ«r tĂ« gjitha bibliotekat e jashtme — ato duhet tĂ« jenĂ« tĂ« instaluara nĂ« makinat me planifikuesin dhe punĂ«torĂ«t.

Tani është e thjeshtë:

$ docker-compose up --scale worker=3

Pasi të ngritët gjithçka, mund të shikoni ndërfaqet e uebit:

Konceptet themelore

Nëse nuk kuptoni asgjë nga këto «dag» të gjitha, ja një fjalor të shkurtër:

  • Planifikuesi — personi kryesor nĂ« Airflow, qĂ« kontrollon qĂ« robotĂ«t tĂ« punojnĂ«, jo njeriu: monitoron orarin, pĂ«rditĂ«son dagĂ«t, nis detyrat.

    NĂ« fakt, nĂ« versionet e vjetra, ai kishte probleme me memorjen (jo, jo amnezi, por rrjedhje) dhe nĂ« konfigurime kishte njĂ« parametrin legaciy run_duration — intervali i ripĂ«rsĂ«ritjes sĂ« tij. Por tani gjithçka Ă«shtĂ« nĂ« rregull.

  • DAG (i njohur si «dag») — «graf me drejtime jo-ciklike», por kjo pĂ«rkufizim nuk thotĂ« shumĂ« pĂ«r shumicĂ«n, dhe nĂ« thelb Ă«shtĂ« njĂ« enĂ« pĂ«r detyrat qĂ« ndĂ«rveprojnĂ« me njĂ«ra-tjetrĂ«n (shih mĂ« poshtĂ«) ose njĂ« analog PaketĂ« nĂ« SSIS dhe Fleksibilitet nĂ« Informatica.

    Përveç dagëve, mund të ketë edhe subdagë, por ndoshta nuk do të arrijmë të merremi me ta.

  • Kalim i DAG — njĂ« dag i inicializuar, tĂ« cilit i Ă«shtĂ« caktuar execution_date. Kalimet e njĂ« dag-u mund tĂ« punojnĂ« paralelisht (nĂ«se natyrisht i keni bĂ«rĂ« detyrat tuaja idempotente).
  • Operator — ato janĂ« copa kodi qĂ« janĂ« pĂ«rgjegjĂ«se pĂ«r kryerjen e ndonjĂ« veprimi tĂ« caktuar. Ka tre lloje operatorĂ«sh:
    • action, siç Ă«shtĂ« operatori ynĂ« i dashur PythonOperator, i cili Ă«shtĂ« nĂ« gjendje tĂ« ekzekutojĂ« çdo kod (tĂ« vlefshĂ«m) Python;
    • transfer, qĂ« transferojnĂ« tĂ« dhĂ«na nga njĂ« vend nĂ« njĂ« tjetĂ«r, e thĂ«nĂ« ndryshe, MsSqlToHiveTransfer;
    • sensor do tĂ« lejojĂ« tĂ« reagoni ose ngadalĂ«sojĂ« ekzekutimin e dag-ut deri nĂ« ndodhin e njĂ« ngjarjeje. HttpSensor mund tĂ« thĂ«rrasĂ« endpoint-in e specified, dhe kur pret pĂ«rgjigjen e nevojshme, nis transferimin GoogleCloudStorageToS3Operator. NjĂ« mendje kurioze do tĂ« pyesĂ«: «pĂ«rse? Sepse Ă«shtĂ« e mundur tĂ« bĂ«sh pĂ«rsĂ«ritje direkt nĂ« operator!» Pastaj, pĂ«r tĂ« mos mbushur pulin e detyrave me operatorĂ« tĂ« ngrirĂ«. Sensori aktivizohet, kontrollon dhe ndalet deri nĂ« pĂ«rpjekjen tjetĂ«r.
  • DetyrĂ« — operatorĂ«t e shpallur, pa marrĂ« parasysh llojin e tyre dhe tĂ« lidhur me dag-u, ngrihen nĂ« gradĂ«n e detyrĂ«s.
  • InstancĂ« detyre — kur plani gjeneral vendosi se detyrat Ă«shtĂ« koha tĂ« dĂ«rgohen nĂ« luftarĂ«t ekzekutues (direkt nĂ« vend, nĂ«se pĂ«rdorim LocalExecutor apo nĂ« njĂ« nodĂ« tĂ« largĂ«t nĂ« rastin e CeleryExecutor), ai u caktton kontekst (dmth, njĂ« set variables — parametrat e ekzekutimit), zhvillon shabllonet e komandave ose kĂ«rkesave dhe i vendos ato nĂ« pool.

Generojmë detyra

Së pari, do të shënojmë skemën e përgjithshme të dag-ut tonë, dhe më pas do të zhytim gjithnjë e më thellë në detaje, sepse ne po aplikojmë disa zgjidhje tejet të veçanta.

Kështu, në formën më të thjeshtë, një dag i tillë do të duket kështu:

from datetime 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)

Le të kuptojmë:

  • PĂ«r sĂ«pari importojmĂ« bibliotekat e nevojshme dhe diçka tjetĂ«r;
  • sql_server_ds — Ă«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 ËshtĂ« shpallja e dag-ut tonĂ«, e cila duhet tĂ« jetĂ« nĂ« globals(), ndryshe Airflow nuk do ta gjejĂ«. Po ashtu, duhet tĂ« informojmĂ« dag-un:
    • si e quajmĂ« orders Ky do tĂ« jetĂ« emri qĂ« do tĂ« shfaqet nĂ« ndĂ«rfaqen e uebit,
    • se ai do tĂ« fillojĂ« punĂ«n nga mesnata e 8 korrikut,
    • dhe duhet tĂ« ekzekutohet, pĂ«rafĂ«rsisht çdo 6 orĂ« (pĂ«r djemtĂ« cool kĂ«tu nĂ« vend tĂ« timedelta() lejohet cron-string 0 0 0/6 ? * * *, pĂ«r ata qĂ« nuk janĂ« aq cool — njĂ« shprehje si @daily);
  • workflow() do tĂ« bĂ«jĂ« punĂ«n kryesore, por jo tani. Tani thjesht do ta hedhim kontekstin tonĂ« nĂ« log.
  • Tani magjia e thjeshtĂ« e krijimit tĂ« detyrave:
    • kalojmĂ« nĂ«pĂ«r burimet tona;
    • initializo PythonOperator, e cila do tĂ« ekzekutojĂ« zbrazĂ«tinĂ« tonĂ« workflow(). Mos harroni tĂ« jepni njĂ« emĂ«r unik (brenda dag-ut) pĂ«r detyrĂ«n dhe tĂ« lidhni vetĂ« dag-un. Flag provide_context do tĂ« çojĂ« nĂ« funksion disa argumente shtesĂ«, tĂ« cilat ne do t'i mbledhim me kujdes me **context.

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

  • njĂ« dag tĂ« ri nĂ« ndĂ«rfaqen e uebit,
  • njĂ«qind e pesĂ«dhjetĂ« detyra, tĂ« cilat do tĂ« ekzekutohen paralelisht (nĂ«se konfigurimet e Airflow, Celery dhe fuqisĂ« sĂ« serverĂ«ve lejojnĂ«).

Më saktësisht, pothuajse kemi arritur.

Apache Airflow: e bëjmë ETL më të thjeshtë
Kush do të vendosë varësitë?

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

Tani, le të fillojmë:

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

KatrorĂ«t gri — instancat e detyrĂ«s, tĂ« pĂ«rpunuara nga planifikuesi.

Pak po presim, detyrat po kapen nga punëtorët:

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

GjelbĂ«rt, natyrisht, - janĂ« tĂ« pĂ«rfunduara me sukses. TĂ« kuq — jo aq me sukses.

NĂ« fakt, nĂ« prodhimin tonĂ« nuk ka asnjĂ« dosje . /dags, e cila sinkronizohet midis makinave — tĂ« gjitha dag-tĂ« ndodhen nĂ« git nĂ« Gitlab-in tonĂ«, dhe CI Gitlab vendos azhurnimet nĂ« makina kur Ă«shtĂ« bĂ«rĂ« bashkimi nĂ« master.

Pak për Flower

NdĂ«rsa punĂ«torĂ«t po pĂ«rfundjnĂ« detyrat tona zbrazĂ«ta, le tĂ« kujtojmĂ« njĂ« instrument tjetĂ«r qĂ« mund tĂ« na tregojĂ« diçka — Flower.

Faqja e parë me informacion të përmbledhur për nodet-punëtorë:

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

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

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

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

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

Faqja mĂ« e ndritshme — me grafikĂ«t e gjendjes sĂ« detyrave dhe kohĂ«s sĂ« ekzekutimit tĂ« tyre:

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

Ngarkojmë atë që nuk është ngarkuar

Pra, të gjitha detyrat u përfunduan, mund të marrim të lënduarit.

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

Dhe ka pasur mjaft tĂ« lĂ«nduar — pĂ«r arsye tĂ« ndryshme. NĂ« rastin e pĂ«rdorimit tĂ« duhur tĂ« Airflow, kĂ«to katrorĂ« flasin pĂ«r faktin se tĂ« dhĂ«nat sigurisht qĂ« nuk arritĂ«n.

Duhet të shikojmë logun dhe të rindezim instancat e taskëve që kanë rënë.

Duke klikuar në çdo katror, do të shohim veprimet e disponueshme:

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

Mund të marrim dhe të bëjmë Clear për atë që ka rënë. Kjo do të thotë, harrojmë se diçka ka shkuar keq dhe ajo instancë e njëjtë do të kthehet në planifikues.

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

E qartĂ«, qĂ« tĂ« veprosh kĂ«shtu me tĂ« gjitha katrorĂ«t e kuq nuk Ă«shtĂ« shumĂ« human — nuk Ă«shtĂ« kjo ajo çfarĂ« presim nga Airflow. Natyrisht, ne kemi armĂ«n tonĂ« tĂ« shkatĂ«rrimit nĂ« masĂ«: Browse/Task Instances

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

Do të zgjedhim të gjitha njëherësh dhe do të zbresim duke shtypur pikën e duhur:

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

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

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

Koneksionet, huki dhe variabla të tjera

ËshtĂ« koha 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 të përditësuara',
    html_content=dedent("""Gospodë të mirë, 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, zgjohet, ne kemi rënë {{ 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]

Ați făcut vreodată un raport de actualizare? Este din nou aceeași poveste: avem o listă de surse din care să extragem datele; avem o listă Ăźn care să le așezăm; să nu uităm să semnalizăm cĂąnd totul s-a ĂźntĂąmplat sau s-a stricat (păi, asta nu e despre noi, nu?).

Haideți să trecem din nou prin fișier și să vedem lucrurile noi neclare:

  • from commons.operators import TelegramBotSendMessage — nimic nu ne Ăźmpiedică să facem propriile noastre operatore, iar noi am profitat de acest lucru făcĂąnd o mică Ăźnfășurare pentru a trimite mesaje Ăźn Descheiat. (Vom vorbi mai jos despre acest operator);
  • default_args={} — DAG-ul poate oferi aceleași argumente tuturor operatoarelor sale;
  • to='{{ var.value.all_the_kings_men }}' — cĂąmp to nu va fi hardcodat, ci generat dinamic cu ajutorul Jinja și al unei variabile cu lista de emailuri pe care am așezat-o cu grijă Ăźn Admin/Variables;
  • trigger_rule=TriggerRule.ALL_SUCCESS — condiția de declanșare a operatorului. În cazul nostru, mesajul va ajunge la șefi doar dacă toate dependențele au rulat cu succes;
  • tg_bot_conn_id='tg_main' — argumente conn_id acceptă identificatorii conexiunilor pe care le creăm Ăźn Admin/Connections;
  • trigger_rule=TriggerRule.ONE_FAILED — mesajele Ăźn Telegram vor pleca doar Ăźn cazul Ăźn care există sarcini căzute;
  • task_concurrency=1 — interzicem lansarea simultană a mai multor instanțe de sarcină pentru aceeași sarcină. În caz contrar, vom obține lansarea simultană a mai multor VerticaOperator (care se uită la aceeași tabelă);
  • report_update >> [email, tg] — totul VerticaOperator se va aduna Ăźn trimiterea unui email și a unui mesaj, iată așa:
    Apache Airflow: e bëjmë ETL më të thjeshtë

    Dar, deoarece condțiile de declanșare a operatoarelor de notificare sunt diferite, va funcționa doar unul. În Tree View, totul arată puțin mai puțin clar:
    Apache Airflow: e bëjmë ETL më të thjeshtë

Voi spune cĂąteva cuvinte despre macro-uri și prietenii lor — variabile.

Macro-urile sunt marcaje Jinja care pot introduce diferite informații utile Ăźn argumentele operatoarelor. De exemplu, astfel:

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

{{ ds }} se va desfășura Ăźn conținutul variabilei de context execution_date Ăźn format YYYY-MM-DD: 2020-07-14. Cel mai plăcut este că variabilele de context sunt legate de un anumit instantaneu de sarcină (cadrul din Tree View), iar la repornirea plasează marcajele Ăźn aceleași valori.

Valorile atribuite pot fi vizualizate folosind butonul Rendered pe fiecare instanță de sarcină. Iată cum arată sarcina cu trimiterea emailului:

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

Iată cum arată sarcina cu trimiterea mesajului:

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

Lista e plotave të integruara për versionin më të fundit është këtu: Referenca për Makro

Për më tepër, me ndihmën e plugjineve, ne mund të shpallim makro të tijat, por kjo është një histori krejt tjetër.

Përveç gjërave të paracaktuara, ne mund të zëvendësojmë vlerat e variablave tanë (më lart në kod kam bërë tashmë këtë). Le të krijojmë disa: Admin/Variables Pika të shumta:

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

Gjithçka është gati për tu përdorur:

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

Një vlerë mund të jetë scalare, por mund të jetë dhe JSON. Në rastin e JSON-it:

bot_config

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

thjesht përdorim rrugën për çelësin përkatës: {{ var.json.bot_config.bot.token }}.

Do të them një fjalë të vetme dhe do të tregoj një ekran për lidhet. Këtu gjithçka është elementare: në faqen Admin/Connections krijojmë lidhjen, vendosim atje emrat tanë të përdoruesve/faqet dhe parametra më specifik. Ja si:

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

FjalĂ«kalimet mund tĂ« jenĂ« tĂ« kriptuara (mĂ« me imtĂ«si se nĂ« variantin e paracaktuar), ose mund tĂ« mos specifikohet lloji i lidhjes (siç bĂ«ra unĂ« pĂ«r tg_main) — çështja Ă«shtĂ« se lista e llojeve Ă«shtĂ« e ngulitur nĂ« modelet Airflow dhe nuk mund tĂ« zgjerohet pa u ndaluar nĂ« burimet (nĂ«se ndonjĂ«herĂ« kam harruar diçka — ju lutem mĂ« korrigjoni), por Ă«shtĂ« e thjeshtĂ« tĂ« marrim kredencialet vetĂ«m me emrin.

Dhe gjithashtu është e mundur të bësh disa lidhje me të njëjtin emër: në një rast të tillë, metoda BaseHook.get_connection(), e cila na sjell lidhjet sipas emrit, do të kthejë një të rastit nga disa të ngjashëm (do të ishte më logjike të bënte Round Robin, por këtë do ta lëmë në përgjegjësinë e zhvilluesve të Airflow).

Variablat dhe Lidhjet janë padyshim mjete të shkëlqyera, por është e rëndësishme të mos humbasësh ekuilibrin: cilat pjesë të proceseve të tua i ruani realisht në kod, dhe cilat i besoni ruajtjes nga Airflow. Nga njëra anë, ndërrimi i shpejtë i vlerës, për shembull, e-maili i dërguesit, mund të jetë i përshtatshëm përmes UI. Nga ana tjetër, kjo është në fund të fundit një kthim në përdorimin e miut, nga të cilin ne (unë) dëshironim të ishim të lirë.

Puna me lidhjet është një nga detyrat hukut. Në të vërtetë hook-et e Airflow janë pika lidhjeje me shërbime dhe biblioteka të jashtme. Për shembull, JiraHook do të hapë për ne një klient për ndërveprim me Jira (mund të lëvizim detyra andej këndej), dhe me ndihmën e SambaHook mund të dërgoni një skedar lokal në pikën smb.Dhe ne po afrohemi për të parë se si është bërë

Shpjegojmë operatorin personalizuar

TelegramBotSendMessage commons/operators.py

Kodi me vetë operatorin: me operatorin vetë:

nga importimi i Union

nga airflow.operators import BaseOperator

nga commons.hooks import TelegramBotHook, TelegramBot

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

    Shembuj:
        >>> 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 }} ka dështuar :(',
        ...     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, si gjithçka tjetër në Airflow, gjërat janë shumë të thjeshta:

  • TrashĂ«guam nga BaseOperator, e cila implementon shumĂ« gjĂ«ra specifike pĂ«r Airflow (shikoni pĂ«r mĂ« shumĂ« kur tĂ« keni kohĂ«)
  • Deklaruam fushat template_fields, nĂ« tĂ« cilat Jinja do tĂ« kĂ«rkojĂ« makros pĂ«r pĂ«rpunim.
  • Organizuam argumentet e duhura pĂ«r __init__(), vendosĂ«m parametra tĂ« paracaktuar, kur Ă«shtĂ« nevoja.
  • Nuk harrojmĂ« as pĂ«r inicializimin e prindit.
  • HapĂ«m hooking pĂ«rkatĂ«s TelegramBotHook, morĂ«m objektin-klient prej tij.
  • E pĂ«rcaktuam (override) metodĂ«n BaseOperator.execute(), e cila do tĂ« thirret nga Airflow kur tĂ« jetĂ« koha pĂ«r tĂ« ekzekutuar operatorin — aty e realizojmĂ« veprimin kryesor, pa harruar tĂ« regjistrohemi. (Regjistrohemi, pĂ«r fat tĂ« keq, pikĂ«risht nĂ« stdout dhe stderr — Airflow do ta kapĂ«, do ta mbajĂ« bukur, do ta rregullojĂ« dhe do ta vendosĂ« atje ku duhet.)

Le të shohim se çfarë kemi në commons/hooks.py. Pjesa e parë e skedarit, me hooking vetë:

nga importimi i Union

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

klasa TelegramBotHook(BaseHook):
    """Huku i API të Telegram Bot

    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

Nuk e di se çfarë mund të shpjegoj këtu, vetëm do të theksoj momentet e rëndësishme:

  • Ne trashĂ«gojmĂ«, mendojmĂ« pĂ«r argumentet — nĂ« shumicĂ«n e rasteve do tĂ« ketĂ« vetĂ«m njĂ«: conn_id;
  • E pĂ«rcaktojmĂ« metodat standarde: unĂ« u kufizova nĂ« get_conn(), nĂ« tĂ« cilin marr parametrat e lidhjes sipas emrit dhe thjesht nxjerr seksionin extra (kjo Ă«shtĂ« fusha pĂ«r JSON), nĂ« tĂ« cilĂ«n unĂ« (sipas udhĂ«zimeve tĂ« mia!) vendosa token e Telegram-botit: {"bot_token": "YOuRAwEsomeBOtToKen"}.
  • Po krijoj njĂ« ekzemplar tĂ« TelegramBot, duke i dhĂ«nĂ« atij tani njĂ« token tĂ« caktuar.

Kjo është gjithçka. Mund të marr klientin nga hooku me anë të TelegramBotHook().clent ose TelegramBotHook().get_conn().

Dhe pjesa e dytë e skedarit, në të cilën unë bëj një mikro-kapsulë për Telegram REST API, në mënyrë që të mos sjell të njëjtin python-telegram-bot për një metodë të vetme sendMessage.

class TelegramBot:
    """Përfaqësuesi i API-së së Telegram Bot

    Shembuj:
        >>> TelegramBot('YOuRAwEsomeBOtToKen', '@myprettydebugchat').send_message('Përshëndetje, dashuri')
        >>> TelegramBot('YOuRAwEsomeBOtToKen').send_message('Përshëndetje, dashuri', 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))

Rruga e duhur Ă«shtĂ« tĂ« bashkosh gjithçka kĂ«tĂ«: commons/operators.py, TelegramBotHook, TelegramBot — nĂ« njĂ« plugin, ta vendosĂ«sh nĂ« njĂ« depo publike, dhe ta ofrosh nĂ« Open Source.

NdĂ«rsa po studiojmĂ« tĂ« gjitha kĂ«to, azhurnimet e raportit tonĂ« kishin arritur tĂ« mblidhen dhe tĂ« mĂ« dĂ«rgojnĂ« njĂ« mesazh gabimi nĂ« kanal. Do shkoj tĂ« kontrolloj se çfarĂ« ka ndodhur pĂ«rsĂ«ri


Apache Airflow: e bëjmë ETL më të thjeshtë
Në DAG-un tonë diçka ka shkuar keq! E ndonjëherë kjo ishte ajo që prisnim? Pikërisht!

A do të derdhësh?

A besoni se unë kam humbur diçka? Duket se propozova të transferoj të dhënat nga SQL Server në Vertica, e këtu kam devijuar nga tema, i paaftë!

Kjo veprim ishte e qëllimshme, thjesht isha i detyruar t'ju shpjegoja disa terminologji. Tani mund të vazhdojmë.

Plani ynë ishte kështu:

  1. Të bëjmë një DAG
  2. Të gjenerojmë detyra
  3. Të shohim se si është gjithçka e bukur
  4. Të caktojmë numrat e sesioneve për ngarkimet
  5. Të marrim të dhënat nga SQL Server
  6. Të vendosim të dhënat në Vertica
  7. TĂ« mbledhim statistika

Prandaj, për ta nisur 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

Aty jemi duke ngritur:

  • Vertica si host dwh me cilĂ«simet mĂ« tĂ« zakonshme,
  • tre instance SQL Server,
  • mbushim bazat nĂ« fund me disa tĂ« dhĂ«na (asnjĂ«herĂ« mos shikoni nĂ« mssql_init.py!)

E aktivizojmë gjithçka me një komandë disi më të komplikuar se herën e kaluar:

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

ÇfarĂ« ka gjeneruar gjeneratori ynĂ« tĂ« mrekullueshĂ«m, mund ta shohim, duke pĂ«rdorur opsionin Data Profiling/Ad Hoc Query:

Apache Airflow: e bëjmë ETL më të thjeshtë
E rëndësishme, mos e tregoni këtë analistëve

Nuk do të ndalem në detaje mbi seancat ETL unë nuk do, aty është krejtësisht e thjeshtë: krijojmë bazën, në të një tavëll, e mbështjellim gjithçka me menaxherin e kontekstit, dhe tani e bëjmë kështu:

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

    # Load workflow
    ...

    session.successful = True
    session.loaded_rows = 15

session.py

nga sys import stderr

klasa Sesioni:
    """Seanca ETL workflow

    Shembuj:
        me Sesioni(emri_detyrës) si seancë:
            print(seancë.id)
            seancë.suksesshëm = True
            seancë.rreshtat_e_ngarkuar = 15
            seancë.koment = 'Shumë mirë'
    """

    def __init__(self, lidhja, emri_detyrës):
        vetë.lidhja = lidhja
        vetë.lidhja.autocommit = True

        vetë._emri_detyrës = emri_detyrës
        vetë._id = None

        vetë.rreshtat_e_ngarkuar = None
        vetë.suksesshëm = None
        vetë.koment = None

    def __enter__(self):
        kthehu vetë.hap()

    def __exit__(self, tipi_exc, vlera_exc, tb_exc):
        nëse any(tipi_exc, vlera_exc, tb_exc):
            vetë.suksesshëm = False
            vetë.koment = f'{tipi_exc}: {vlera_exc}n{tb_exc}'
            print(tipi_exc, vlera_exc, tb_exc, file=stderr)
        vetë.mbyll()

    def __repr__(self):
        kthehu (f'')

    @property
    def emri_detyrës(self):
        kthehu vetë._emri_detyrës

    @property
    def id(self):
        kthehu vetë._id

    def _ekzekuto(self, pyetje, *args):
        me vetë.lidhja.cursor() si kursori:
            kursori.execute(pyetje, args)
            kthehu kursori.fetchone()[0]

    def _krijo(self):
        pyetje = """
            KRIJO TABELËN NËSE NUK EKZISTON sesione (
                id          SERIAL       JO NULL KRYESORE,
                emri_detyrës   VARCHAR(200) JO NULL,

                filluar     TIMESTAMPTZ  JO NULL DEFAULT current_timestamp,
                përfunduar    TIMESTAMPTZ           DEFAULT current_timestamp,
                sukses      BOOL,

                rreshtat_e_ngarkuar INT,
                koment     VARCHAR(500)
            );
            """
        vetë._ekzekuto(pyetje)

    def hap(self):
        pyetje = """
            SHTO NË sesione (emri_detyrĂ«s, pĂ«rfunduar)
            VLERAT (%s, NULL)
            KTHIM ID;
            """
        vetë._id = vetë._ekzekuto(pyetje, vetë.emri_detyrës)
        print(vetë, 'hapur')
        kthehu vetë

    def mbyll(self):
        nëse jo vetë._id:
            ngrije SesionClosedError('Seanca nuk është e hapur')
        pyetje = """
            UPDATE sesione
            SHTO
                përfunduar = DEFAULT,
                sukses      = %s,
                rreshtat_e_ngarkuar = %s,
                koment     = %s
            KU
                id = %s
            KTHIM ID;
            """
        vetë._ekzekuto(pyetje, vetë.suksesshëm, vetë.rreshtat_e_ngarkuar,
                      vetë.koment, vetë.id)
        print(vetë, 'mbyllur',
              ', sukses: ', vetë.suksesshëm,
              ', Ngarkuar: ', vetë.rreshtat_e_ngarkuar,
              ', koment:', vetë.koment)

klasa SesionError(Exception):
    kaloni

klasa SesionClosedError(SesionError):
    kaloni

Ka ardhur koha të marrim të dhënat tona nga gjysma e qindrave të tabelave tona. Le t'i marrim ato me disa rreshta të thjeshtë:

source_conn = MsSqlHook(mssql_conn_id=src_conn_id, schema=src_schema).get_conn()

pyetje = f"""
    ZGJIDH
        id, koha_e_fillimit, koha_e_përfundimit, tipi, të dhënat
    NGA dbo.Orders
    KU
        CONVERT(DATE, koha_e_fillimit) = '{dt}'
    """

df = pd.read_sql_query(pyetje, source_conn)
  1. Me anë të hook ne do të marrim nga Airflow pymssql-konnektimin
  2. Në pyetje do të vendosim një kufizim në formën e datës - në funksion, do ta sjellë nëpërmjet shablonit.
  3. Ne e japim pyetjen tonĂ« pandas, e cila do tĂ« nxjerrĂ« pĂ«r ne MundĂ«sia e qasjes nĂ« tĂ« dhĂ«na nĂ« formatin çelĂ«s-vlerĂ« ose grupe kolonesh ( — do na nevojitet mĂ« pas.

Unë përdor zëvendësimin {dt} në vend të parametrave të pyetjes %s jo sepse unë jam një Buratino i keq, por sepse pandas nuk mund ta përballojë pymssql dhe ia jep të fundit params: Lista, ndonëse ai dëshiron shumë tuple.
Gjithashtu, vini re se zhvilluesi pymssql vendosi të mos e mbështesë më, dhe është koha për t'u zhvendosur në pyodbc.

Le të shohim se me çfarë Airflow i ka mbushur argumentet funksioneve tona:

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

Nëse nuk ka të dhëna, atëherë nuk ka kuptim të vazhdojmë. Por është e çuditshme të mendohet se ngarkimi ishte i suksesshëm. Por kjo gjithashtu nuk është gabim. A-a-a, çfarë të bëjmë?! Ja çfarë:

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

AirflowSkipException do t'i thotë Airflow, që në fakt nuk ka gabim, dhe ne do ta kalojmë detyrën. Në ndërfaqe, nuk do të ketë katror të gjelbër e as të kuq, por një ngjyrë pink.

Le të hedhim për të dhënat tona disa kolona:

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

Pikërisht:

  • Baza e tĂ« dhĂ«nave nga e cila morĂ«m porositĂ«,
  • Identifikuesi i sesionit tonĂ« tĂ« ngarkimit (do tĂ« jetĂ« ndryshe pĂ«r çdo detyrĂ«),
  • Hashi nga burimi dhe identifikuesi i porosisĂ« — pĂ«r tĂ« pasur njĂ« identifikues unik tĂ« porosisĂ« nĂ« bazĂ«n pĂ«rfundimtare (ku tĂ« gjitha do tĂ« bashkohen nĂ« njĂ« tabelĂ«).

Kanë mbetur vetëm dy hapa: ngarkoni gjithçka në Vertica. Dhe, ë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 po krijojmë një marrës të posaçëm StringIO.
  2. pandas do ta grumbullojë në të tonë Mundësia e qasjes në të dhëna në formatin çelës-vlerë ose grupe kolonesh ( në formën e CSV-rreshta.
  3. Do të hapim një lidhje me Vertica-n tonë të dashur përmes hooks.
  4. Dhe tani me ndihmën e copy() do t'i dërgojmë të dhënat tona drejtpërdrejt në Vertica!

Do të marrim nga driver-i se sa rreshta janë ngarkuar, dhe do të informojmë menaxherin e sesionit se gjithçka është në rregull:

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

Kjo është gjithçka.

Në prodhim ne krijojmë tabelën target manualisht. Këtu, unë e lejoj veten një automatik të vogël:

create_schema_query = f'KRIJO SKEMË NËSE NUK EKSPLOZOJ {target_schema};'
create_table_query = f"""
    KRIJO TABELË NËSE NUK EKSPLOZOJ {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)

Me ndihmën e VerticaOperator() krijoj skemën e bazës së të dhënave dhe tabelën (nëse ato nuk ekzistojnë, natyrisht). E rëndësishme është të vendosni saktë 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

Të shohim përfundimet

— Ja, — tha miu, — nuk Ă«shtĂ« e vĂ«rtetĂ« se tani
A e ke verifikuar se në pyll unë jam kafsha më e frikshme?

Julia Donaldson, "Gruffalo"

Mendoj se po tĂ« bĂ«nim me kolegĂ«t e mi njĂ« garĂ«: kush e krijon dhe e nis mĂ« shpejt njĂ« proces ETL nga fillimi: ata me SSIS dhe miun e tyre dhe unĂ« me Airflow
 Dhe pastaj do tĂ« krahasonim lehtĂ«sinĂ« e mbajtjes
 Euh, mendoj se do tĂ« pajtoheshit qĂ« unĂ« do t'i kaloj nĂ« çdo aspekt!

Nëse flasim pak më seriozisht, Apache Airflow - për shkak se përshkrimi i proceseve bëhet në formën e kodit programor - e bëri punën time shumë më të lehtë dhe më të këndshme.

Papërjashtim, zgjerueshmëria e pafundme: si në planin e plugin-eve, ashtu dhe predispozita për shkallëzim - ju jep mundësinë të aplikoni Airflow praktikisht në çdo fushë: qoftë në ciklin e plotë të mbledhjes, përgatitjes dhe përpunimit të të dhënave, qoftë në nisjen e raketave (në Mars, sigurisht).

Pjesa përfundimtare, informuese dhe referuese

Gramat për të cilat ne i mbledhëm për ju

  • start_date. Po, ky Ă«shtĂ« njĂ« meme lokal tani. NĂ«pĂ«rmjet argumentit kryesor tĂ« DAG-ut start_date kalojnĂ« tĂ« gjithĂ«. NĂ« mĂ«nyrĂ« tĂ« shkurtĂ«r, nĂ«se e vendosni nĂ« start_date datĂ«n aktuale, dhe nĂ« schedule_interval nĂ« njĂ« ditĂ«, atĂ«herĂ« DAG-u do tĂ« nisĂ« nesĂ«r jo mĂ« herĂ«t.
    start_date = datetime(2020, 7, 7, 0, 1, 2)

    Dhe më shumë asnjë problem.

    Kjo është gjithashtu e lidhur me një gabim tjetër ekzekutimi: Task is missing the start_date parameter, i cili shpesh tregon se keni harruar të lidhni operatorin me DAG-un.

  • E gjithĂ« nĂ« njĂ« makinĂ«. Po, dhe bazat (e Airflow-it dhe mbulesĂ«s sonĂ«), serveri i web-it, planifikuesi dhe punĂ«torĂ«t. Dhe madje ka funksionuar. Por me kalimin e kohĂ«s, numri i detyrave nĂ« shĂ«rbime u rrit, dhe kur PostgreSQL filloi tĂ« kthejĂ« pĂ«rgjigje pĂ«r indeksin pĂ«r 20 ms nĂ« vend tĂ« 5 ms, ne e morĂ«m dhe e larguam.
  • LocalExecutor. Po, ne jemi ende nĂ« tĂ«, dhe tani jemi afĂ«r skajit. LocalExecutor na ka mjaftuar deri tani, por tani ka ardhur koha tĂ« zgjerohemi me tĂ« paktĂ«n njĂ« punĂ«tor, dhe do tĂ« duhet tĂ« punojmĂ« pĂ«r tĂ« kaluar nĂ« CeleryExecutor. Dhe duke marrĂ« parasysh qĂ« me tĂ« mund tĂ« punoni edhe nĂ« njĂ« makinĂ«, nuk ka asgjĂ« qĂ« ndalon pĂ«rdorimin e Celery as nĂ« serverin qĂ« "natyrisht, kurrĂ« nuk do tĂ« shkojĂ« nĂ« prodhim, e betohem!"
  • Mos pĂ«rdorimi i mjeteve tĂ« integruara:
    • Connections pĂ«r ruajtjen e kredencialeve tĂ« shĂ«rbimeve,
    • SLA Misses pĂ«r t'iu pĂ«rgjigjur detyrave qĂ« nuk pĂ«rfunduan me kohĂ«,
    • XCom pĂ«r exchange tĂ« metadonnĂ©es (thashĂ« metagtĂ« dhĂ«nave!) midis detyrave tĂ« DAG-ut.
  • KeqpĂ«rdorimi i postĂ«s. ÇfarĂ« mund tĂ« themi kĂ«tu? JanĂ« konfiguruar njoftime pĂ«r tĂ« gjitha pĂ«rsĂ«ritjet e detyrave tĂ« rĂ«na. Tani nĂ« Gmail tim tĂ« punĂ«s kam >90k email-a nga Airflow, dhe ndĂ«rfaqja web e postĂ«s refuzon tĂ« marrĂ« dhe fshijĂ« mĂ« shumĂ« se 100 njĂ«si njĂ«herĂ«sh.

Më shumë gurë nën ujë: Apache Airflow Pitfails

Mjetet për automatizim edhe më të madh

Për të punuar edhe më shumë me mendje sesa me duar, Airflow na ka përgatitur këto:

  • REST API — akoma ka statusin Eksperimental, qĂ« nuk e pengon tĂ« punojĂ«. Me tĂ« Ă«shtĂ« e mundur jo vetĂ«m tĂ« marrim informacion mbi DAG-Ă«t dhe detyrat, por tĂ« ndalojmĂ«/nisim njĂ« DAG, tĂ« krijojmĂ« njĂ« DAG Run apo njĂ« pool.
  • CLI — pĂ«rmes dhomĂ«s sĂ« komandave, shumĂ« mjete janĂ« tĂ« disponueshme, tĂ« cilat jo vetĂ«m qĂ« janĂ« tĂ« papraktyse pĂ«r tu pĂ«rdorur pĂ«rmes WebUI, por nĂ« fakt nuk ekzistojnĂ« aspak. PĂ«r shembull:
    • backfill nevojitet pĂ«r tĂ« rinisur instancat e detyrave.
      Për shembull, vijnë analistët, thonë: "E ke, shok, një problem në të dhënat nga 1 deri në 13 janar! Niset nga fillimi!". Dhe ti bën:
      airflow backfill -s '2020-01-01' -e '2020-01-13' orders
    • Kujdesi i bazĂ«s: initdb, resetdb, upgradedb, checkdb.
    • run, i cili lejon tĂ« niset njĂ« instancĂ« detyre, duke injoruar tĂ« gjitha varĂ«sitĂ«. MĂ« shumĂ«, Ă«shtĂ« e mundur ta nisim atĂ« pĂ«rmes LocalExecutor, madje edhe nĂ«se ke njĂ« klastĂ«r Celery.
    • PĂ«rshkruan nĂ« mĂ«nyrĂ« tĂ« ngjashme test, por nuk shkruan asgjĂ« nĂ« bazĂ«.
    • connections lejon krijimin nĂ« masĂ« tĂ« lidhjeve nga shell.
  • Python API — njĂ« mĂ«nyrĂ« bastante hardcore pĂ«r tĂ« ndĂ«rvepruar, e cila Ă«shtĂ« e destinuar pĂ«r pluginat, jo pĂ«r tĂ« bĂ«rĂ« punĂ« manuale. Por kush na ndalon tĂ« shkojmĂ« nĂ« /home/airflow/dags, tĂ« nisĂ« ipython dhe tĂ« fillojmĂ« tĂ« çmendemi? Mund, pĂ«r shembull, tĂ« eksportojmĂ« tĂ« gjitha lidhjet me kĂ«tĂ« kod:
    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)
  • Lidhja me bazĂ«n e tĂ« dhĂ«nave tĂ« metadata Airflow. Nuk rekomandoj tĂ« shkruani nĂ« tĂ«, por Ă«shtĂ« shumĂ« mĂ« e shpejtĂ« dhe mĂ« e lehtĂ« tĂ« nxjerrĂ«sh statuset e detyrave pĂ«r metrika tĂ« ndryshme specifike sesa pĂ«rmes ndonjĂ« API.

    TĂ« themi, jo tĂ« gjitha detyrat tona janĂ« idempotente dhe ndonjĂ«herĂ« mund tĂ« bjerĂ« dhe kjo Ă«shtĂ« nĂ« rregull. Por disa grumbuj — kjo Ă«shtĂ« e dyshimtĂ«, dhe duhet tĂ« kontrollojmĂ«.

    Kujdes, SQL!

    ME last_executions SI (
    Zgjidhni
        task_id,
        dag_id,
        execution_date,
        state,
            row_number()
            MBI (
                PARTITION BY task_id, dag_id
                RENDIT
                execution_date DESC) SI rn
    Nga public.task_instance
    KU
        execution_date > tani() - INTERVAL '2' DITË
    ),
    failed SI (
        Zgjidhni
            task_id,
            dag_id,
            execution_date,
            state,
            Rasti Kur rn = row_number() MBI (
                PARTITION BY task_id, dag_id
                RENDIT
                execution_date DESC)
                     ATËHERË E VËRTETË
            SI last_fail_seq
        Nga last_executions
        KU
            state IN ('failed', 'up_for_retry')
    )
    Zgjidhni
        task_id,
        dag_id,
        count(last_fail_seq)                       SI të pasuksesshme,
        count(Rasti Kur last_fail_seq
            DHE state = 'failed' ATËHERË 1 END)       SI dĂ«shtuar,
        count(Rasti Kur last_fail_seq
            DHE state = 'up_for_retry' ATËHERË 1 END) SI up_for_retry
    Nga failed
    GRUPI NGA
        task_id,
        dag_id
    HAVING
        count(last_fail_seq) > 0

Linket

E natyrshme, dhjetë të parët lidhjet nga rezultatet e google janë përmbajtja e dosjes Airflow nga bookmark-at e mia.

Dhe linket e përfshira në artikull:

Burimi: habr.com

Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS đŸ”„ Blini hosting tĂ« besueshĂ«m pĂ«r faqe interneti me mbrojtje nga DDoS, serverĂ« VPS VDS | ProHoster