Strona główna▸Artykuły▸

Apache Airflow i model DAG dla orkiestracji potoków danych

Jak przepływ powietrza wykorzystuje grafy oparte na cyklach nieskończonych (DAG) do harmonogramowania, ponownych prób i uzupełniania danych w potokach danych, oraz dlaczego jawne modelowanie zależności przewyższa łańcuchowe zadania cron razem.

mysimulator teamZaktualizowano — czerwiec 2026≈ 5 min czytania▶ Otwórz symulację

Problem z cronem

Przed pojawieniem się narzędzi do orkiestracji, typowy zespół danych definiował swój 'potok' jako zbiór zadań cron, każde uruchamiane o stałej godzinie i licząc na to, że poprzednie zadanie już zakończyło działanie. Działało to, dopóki nie przestało: zadanie wykonane dwanaście minut po terminie z powodu obciążenia magazynu cicho uszkadzało dane wyjściowe wszystkiego zaplanowanego po nim, a nikt się o tym nie dowiedział do momentu, gdy analityk zapytał, dlaczego liczby na wskaźniku z poprzedniego dnia wyglądały inaczej. Cron koduje tylko czas, nigdy zależność, dlatego jest strukturalnie niezdolny do wyrażania 'zadanie B powinno wykonać się po pomyślnym zakończeniu zadania A, a nie jedynie po ustalonym przesunięciu w czasie.'

Rdzeniową zmianą Airflow było zastąpienie wyzwalania opartego na czasie grafem zależności. Potok jest wyrażany jako Kierowany Graf Niecykliczny, lub DAG: zbiór zadań (węzłów) połączonych krawędziami, które oznaczają 'to zadanie musi zostać zakończone przed rozpoczęciem tego'. Ponieważ graf nie jest cykliczny, nie ma możliwości przypadkowego utworzenia nieskończonej pętli zależności – Airflow weryfikuje to podczas analizy i odmawia zaplanowania DAG-a, który nie jest prawidłowym DAG-iem. Zadaniem planisty staje się wtedy przebywanie grafu: przejście przez graf, znajdowanie zadań, których upstreamowe zależności zakończyły się pomyślnie, i przekazywanie ich do pracownika.

Daty wykonania i idempotentne uruchomienia

Jedną z bardziej subtelnych, ale istotnych decyzji projektowych Airflow jest to, że każdy uruchomienie DAG jest oznaczane datą logicznego wykonania, która reprezentuje okres danych, który ma przetworzyć. Ta data nie jest taka sama jak rzeczywisty czas wykonania uruchomienia. DAG zaplanowany do uruchamiania codziennie o 2:00 dla '2026-07-24' przetwarza dane z 24 dnia, nawet jeśli zostanie uruchomiony późno lub ponownie w trzy dni później – to data logiczna, a nie zegar, określa, na jakich danych operuje zadanie. Ta odseparowanie umożliwia backfilling: jeśli zmienisz transformację i będziesz musiał ponownie przetworzyć ostatnie 90 dni, po prostu wyzwolisz 90 uruchomień DAG, jedno dla każdego historycznego logicznego dnia, a każde z nich zapytuje swój własny dzień's partycję zamiast 'dzisiaj'.

Aby to działało bezpiecznie, zadania muszą być idempotentne – uruchamianie tego samego logicznego dnia dwa razy powinno dawać ten sam wynik, nie podwójnie policzone wiersze. Airflow nie wymusza idempotencji dla Ciebie; jest to umowa, którą musi zawrzeć autor potoku, zwykle poprzez ustawienie każdego zadania tak, aby nadpisywało określoną partycję (`INSERT OVERWRITE PARTITION` lub odpowiednik) zamiast bezwzględnego dodawania. Zespoły, które pomijają tę zasadę, odkrywają koszt w pierwszej kolejności, gdy chwilowe przerwanie zmusza do ponownego uruchomienia i ich tabela przychodów cicho podwaja się.

Operatory, czujniki i abstrakcja zadań

Węzły DAG są budowane z operatorów, które są po prostu parametryzowanymi jednostkami pracy – BashOperator uruchamia polecenie w powłoce, PythonOperator uruchamia wywołanie Pythona, SnowflakeOperator uruchamia zapytanie SQL przeciwko magazynowi danych i tak dalej. Jest to celowo ogólne: Airflow nie obchodzi, co zadanie robi, tylko kiedy powinno ono wykonywać się względem sąsiednich zadań oraz jaka jest jego polityka ponownych prób. Każde zadanie niezależnie może określić liczbę ponownych prób, opóźnienie między nimi (często z wykładniczymi wzrostami), limit czasu, a także pule powiadomień, dzięki czemu awaryjny wywołanie API, które zawodzi raz na 20 próbach, może być automatycznie ponawiane bez powodowania niepowodzenia całej sekwencji ani budzenia nikogo.

Czujniki to specjalny rodzaj operatora, który nie wykonuje pracy, a jedynie czeka na warunek wstępny – FileSensor sprawdza, aż pojawi się plik w chmurze, ExternalTaskSensor czeka, aż zadanie w innym DAG zakończy się, a sensor partycji czeka, aż dane z określonego dnia dotrą do tabeli źródłowej. Czujniki pozwalają na wyrażanie zależności między systemami ('nie uruchamiaj tej transformacji, dopóki nie zostanie napisane dzisiejsze dane przez zadanie w trzech zespołach dalej') bez wpisywania stałego czasu oczekiwania, co jest tą samą kategorią błędu, który promował cron w tamtych czasach. Nowsze wersje Airflow obsługują również 'odroczone' czujniki, które uwalniają swój slot roboczy podczas czekania, dzięki czemu flota potoków oczekujących na wolne systemy upstream nie wyczerpuje planistyka z zasobów.

Scheduler i to, co w rzeczywistości wykonuje

Pod spodem, scheduler Airflow nieustannie analizuje pliki definicji DAG (skrypty Python, które budują obiekty grafu), oceniając, jakie instancje zadań są kwalifikujące się do uruchomienia na podstawie ich zależności i harmonogramu DAG oraz umieszcza je w kolejce. Oddzielny executor decyduje, jak te zadania są faktycznie wykonywane: najprostszy, SequentialExecutor, uruchamia jedno zadanie na raz i nadaje się tylko do lokalnego testowania; CeleryExecutor i KubernetesExecutor rozkładają zadania na grupę pracowników lub tworzą nowy pod dla każdego zadania, odpowiednio – to właśnie w produkcyjnych środowiskach są faktycznie używane. Ta separacja – scheduler decyduje, co ma być uruchomione, a executor decyduje, gdzie – odzwierciedla tę samą zasadę podziel i zwyciężysz, widoczną w podziale ResourceManager versus NodeManager w YARN.

Wizualnie, rendered DAG w interfejsie użytkownika Airflow wygląda jak fabryka: pudełka dla każdego zadania, oznaczona kolorami stanem (w kolejce, uruchomione, zakończone sukcesem, nieudane, nieudane z powodu upstream), połączone strzałkami pokazującymi przepływ zależności, a także widok Gantta, który pokazuje dokładnie, ile czasu zajęło każde zadanie w stosunku do sąsiednich. Ta wizualizacja nie jest przypadkowa – jest to podstawowe narzędzie do debugowania, ponieważ gdy coś się psuje o 3:00 nad ranem, pierwszym pytaniem zawsze jest 'który węzeł w grafie stał się czerwony i jak jego awaria rozeszła się w dół?' a DAG pozwala na odpowiedź na oba te pytania jednym spojrzeniem.

Często zadawane pytania

Czym jest 'DAG' w Airflow?

Oznacza to Graf Kierunkowy Bez Cykli: zbiór zadań połączonych krawędziami o kierunkowej zależności, bez cykli, co gwarantuje zawsze prawidłowy kolejność ich wykonywania.

Dlaczego idempotencja jest tak ważna dla zadań Airflow?

Ponieważ Airflow próbuje ponownie uruchomić nieudane zadania i obsługuje ponowne uruchamianie historycznych dat w celu odzyskiwania danych. Zadanie, które nie jest idempotentne – na przykład takie, które dodaje zamiast nadpisywać partycje – będzie cicho generować duplikaty lub nieprawidłowe dane przy każdym ponownym uruchomieniu lub odzyskaniu.

Jakie jest różnice między harmonizatorem a wykonawcą?

Harmonizator analizuje DAGi i decyduje, które instancje zadań są gotowe do uruchomienia na podstawie zależności i harmonogramu; wykonawca jest oddzielnym komponentem, który określa, jak i gdzie te gotowe zadania są faktycznie wykonywane, lokalnie, poprzez pulę robotników Celery lub jako kontenery Kubernetes.

Jak backfill różni się od normalnego uruchomienia zaplanowanego?

Backfill wywołuje uruchamianie DAGów dla zakresu historycznych dat logicznych zamiast czekać, aż harmonizator do nich dojdzie, pozwalając na ponowne przetwarzanie danych historycznych z zaktualizowaną rurką bez czekania rzeczywistych dni na to, aby harmonizator dogadał się.

Wypróbuj na żywo

Wszystko powyżej działa bezpośrednio w Twojej przeglądarce — otwórz the simulation i zmieniaj parametry podczas działania. Nic nie jest instalowane ani przesyłane na serwer, cały model działa w jednej karcie.

▶ Otwórz symulację the simulation

Co znalazłeś?

Dodaj kroki odtworzenia (opcjonalnie)