ETL w DataFlow

Cykl życia danych i potoków danych

Potok danych na DataFlow AI Platform przebywa długą drogę: jest tworzony w Studiu Projektowym, walidowany, harmonogramowany, kompilowany do DAG, wykonywany zadanie po zadaniu, monitorowany, śledzony pod kątem pochodzenia danych, certyfikowany przez zarządzanie i ostatecznie wycofywany. Ta strona śledzi tę drogę od początku do końca, a także obejmuje cykle życia migracji bazy danych i wydania oprogramowania, które ją wspierają.


Przegląd

Na platformie równolegle przebiegają trzy cykle życia, które przecinają się w przewidywalnych punktach.

Cykl życiaCo przez niego przepływaZarządzany przez
Cykl życia potoku danychDefinicja potoku danych — od wersji roboczej do wycofaniametadata-service + pipeline-engine
Cykl życia uruchomieniaPojedyncze wykonanie potoku danychpipeline-engine, obserwowany przez monitor-service
Cykl życia danychZbiory danych wytwarzane przez potok — skatalogowane, śledzone, zarządzanemetadata-service + lineage-service

Wspierający cykl życia migracji Flyway rozwija schemat bazy danych, a cykl życia wydania oprogramowania dostarcza samą platformę. Pozostała część tej strony omawia każdy z nich po kolei.

  AUTHOR ──► VALIDATE ──► SCHEDULE ──► EXECUTE ──► MONITOR ──► LINEAGE ──► GOVERN ──► RETIRE
   (draft)    (DSL +       (cron /      (DAG run)   (alerts,    (dataset/   (certify,   (archive)
              client       backfill)                cost,SLA)   column)     contracts)
              checks)

Etap 1 — Tworzenie (wersja robocza)

Potok danych rozpoczyna swoje życie jako definicja. Istnieją trzy punkty wejścia, a wszystkie z nich wytwarzają ten sam artefakt: definicję potoku danych DataFlow w formacie YAML.

Punkt wejściaNarzędzieUwagi
Projektowanie wizualneStrona DesignStudio, kanwa React FlowWęzły metodą przeciągnij i upuść; kanwa serializuje się do YAML
Tworzenie z przewodnikiemCreatePipelineWizard (/pipelines/new)Kreator krok po kroku dla typowych kształtów
Import migracjimigration-engineKonwertuje starsze zadanie Informatica / Alteryx / SSIS / DataStage na DataFlow YAML
Instancjonowanie szablonuGaleria PipelineTemplatesRozpoczyna od gotowego szablonu, np. subscriber-360.yaml

DSL w YAML jest parsowany przez dsl/PipelineYamlParser.kt względem dsl/PipelineYamlSchema.kt w pipeline-engine. Tworzenie świadome katalogu jest wspierane, ponieważ metadata-service przechowuje już odkryte schematy — konektory wywołują discoverSchema(), a wynikowe tabele i kolumny są utrwalane w katalogu, dzięki czemu projektant może odwoływać się do prawdziwych zbiorów danych.

Roboczy potok danych i jego wersje są przechowywane jako encje pipelines i pipeline_versions w metadata-service (tworzone przez migracje metadata Flyway V1–V2). Każde zapisanie w Studiu Projektowym tworzy nową wersję; git/PipelineVersionControl.kt i GitRepositoryManager.kt z pipeline-engine podpierają to historią Git, a YamlDiffEngine.kt wytwarza różnice pokazywane w komponencie GitDiffViewer.

Wersjonowanie oparte na Git

Definicje potoków danych są wersjonowane w prawdziwym repozytorium Git zarządzanym przez pipeline-engine. Komponenty VersionHistory i GitDiffViewer w Studiu Projektowym uwidaczniają tę historię; punkty końcowe /api/v1/git udostępniają repozytoria, historię, różnice, wycofanie, gałęzie i scalanie.


Etap 2 — Walidacja

Zanim potok danych będzie mógł zostać uruchomiony, jest walidowany — dwukrotnie.

  1. Po stronie klienta. Usługa SPA services/pipelineValidator.ts sprawdza definicję w przeglądarce w trakcie edycji przez użytkownika, dając natychmiastową informację zwrotną w Studiu Projektowym.
  2. Po stronie serwera. Gdy zażądane jest uruchomienie, validation/PipelineValidator.kt z pipeline-engine waliduje sparsowaną definicję, a validation/ParameterResolver.kt rozwiązuje zmienne wbudowane, zmienne środowiskowe i parametry czasu wykonania.

Walidacja strukturalna kontynuuje się w konstrukcji DAG: dag/DagBuilder.kt buduje ExecutionDAG i odrzuca definicje z cyklami lub wiszącymi zależnościami.

Walidacja kształtu danych to odrębne zagadnienie obsługiwane przez zarządzanie: kontrakty danych (/api/v1/contracts, Flyway metadata V16) deklarują oczekiwany schemat i semantykę zbioru danych, a przepływ pracy migracji udostępnia dedykowaną ValidationSuitePage dla skonwertowanych potoków danych.


Etap 3 — Harmonogramowanie

Zwalidowany potok danych może być uruchamiany na żądanie lub umieszczony w harmonogramie. Harmonogramowanie jest zarządzane przez pakiet scheduler/ z pipeline-engine.

KomponentRola
scheduler/CronScheduler.ktPowtarzające się harmonogramy oparte na cron
scheduler/BackfillManager.ktZadania uzupełniające na partycjach historycznych
scheduler/DependencyChain.ktPorządkowanie zależności między potokami danych
scheduler/RetryPolicy.ktZachowanie ponawiania dla nieudanych uruchomień

Harmonogramy, zadania uzupełniające i zależności są utrwalane jako encje engine_schedules, engine_backfill_jobs / engine_backfill_partitions oraz encje łańcucha zależności (Flyway engine V1 i V3). SchedulerController pod /api/v1/scheduler udostępnia CRUD harmonogramów oraz wstrzymanie i wznowienie, podgląd /next-runs oraz CRUD uzupełnień z wstrzymaniem, wznowieniem i anulowaniem.

Europe/Warsaw domyślnie

Harmonogramy potoków danych domyślnie używają strefy czasowej Europe/Warsaw — platforma jest zbudowana dla Polkomtel, polskiego operatora telekomunikacyjnego, a moduły zgodności RODO / zgodności polskiej zakładają to samo.

Potoki danych mogą być również wyzwalane przez zewnętrzne orkiestratory. Pakiet orchestrator/ dostarcza AirflowClient, AirflowDagGenerator, AutomateNowAdapter oraz WebhookCallbackService, udostępniane przez /api/v1/orchestrator.


Etap 4 — Wykonanie (uruchomienie DAG)

Uruchomienie rozpoczyna się, gdy POST /api/v1/pipelines/{id}/run dociera do ExecutionController z pipeline-engine. execution/PipelineRunner.kt orkiestruje uruchomienie:

  POST /pipelines/{id}/run


  PipelineRunner
    1. PipelineYamlParser   — parse YAML DSL
    2. ParameterResolver    — resolve params / env / runtime vars
    3. PipelineValidator    — validate the definition
    4. DagBuilder           — build ExecutionDAG (cycle / dangling detection)
    5. execute tasks in topological order
         ├─ parallel within each DAG level (fixed thread pool)
         └─ ExecutionContext supports cooperative cancellation
    6. aggregate results

        ├─ ExecutionEventPublisher / PipelineRunLogPublisher  → live events
        ├─ OpenLineageEmitter  → lineage-service /openlineage/events
        └─ on task failure → SelfHealingService

Każde uruchomienie jest zapisywane jako encja engine_execution_runs (z run_id i kolumną JSONB task_states, Flyway engine V4). W przypadku przetwarzania strumieniowego lub obliczeniowo intensywnego silnik zleca zadania zewnętrznym silnikom — flink/FlinkJobSubmitter lub spark/DataprocJobSubmitter.

Samonaprawa

Gdy zadanie zawiedzie, healing/SelfHealingService.kt uruchamia FailureClassifier, aby skategoryzować awarię, i stosuje pasującą strategię z RecoveryStrategies. Próby naprawy są utrwalane (engine_healing_attempts, Flyway engine V8) i uwidaczniane przez HealingController oraz stronę SelfHealingDashboard.

Stany uruchomienia

Uruchomienie przechodzi przez zdefiniowany zestaw stanów.

        ┌──────────┐
        │  QUEUED  │  run requested, awaiting a slot
        └────┬─────┘

        ┌──────────┐
        │ RUNNING  │  tasks executing in topological order
        └────┬─────┘

   ┌─────────┼──────────┬──────────────┐
   ▼         ▼          ▼              ▼
┌────────┐┌────────┐┌──────────┐ ┌──────────┐
│SUCCESS ││ FAILED ││CANCELLED │ │ HEALING  │
└────────┘└───┬────┘└──────────┘ └────┬─────┘
              │  cancellation via      │ recovery strategy
              │  ExecutionContext      │ applied; may
              │                        ▼ re-enter RUNNING
              └──────────────────────────┘
Stan uruchomieniaZnaczenie
QUEUEDUruchomienie zostało zażądane i oczekuje na slot wykonawczy.
RUNNINGDAG jest wykonywany; zadania działają w porządku topologicznym, równolegle w obrębie poziomu.
SUCCESSWszystkie zadania zakończyły się pomyślnie.
FAILEDZadanie zawiodło i nie udało się go odzyskać.
CANCELLEDUruchomienie zostało anulowane kooperacyjnie przez ExecutionContext.
HEALINGSelfHealingService stosuje strategię odzyskiwania; uruchomienie może ponownie wejść w stan RUNNING.

Anulowanie jest udostępniane pod POST /api/v1/runs/{runId}/cancel; status uruchomienia pod GET /api/v1/runs/{runId}.


Etap 5 — Monitorowanie

Po zakończeniu uruchomienia monitor-service pobiera jego metryki i ocenia je. Monitorowanie jest ciągłe, a nie jednorazowym krokiem.

ZagadnienieKomponent monitor-serviceUwidaczniane pod
Status i dzienniki uruchomieniaPipelineRunLogService, PipelineRunLogController/monitor/runs, LogViewerPage
AlertyAlertService, AlertDefinitionService, SseAlertController/monitor/alerts
KosztPipelineCostService, CostAnomalyDetector/monitor/costs
SLA / SLOSlaBurnRateService/monitor/sla
ŚwieżośćFreshnessTracker/monitor/freshness
WydajnośćPerformanceMetricsService/monitor/performance
Bazy odniesienia anomaliiAnomalyBaselineService/monitor/anomalies

Zdarzenia uruchomienia są strumieniowane na żywo do SPA na dwa sposoby: trasa WebSocket (/api/v1/runs/{runId}/stream, obsługiwana przez PipelineStatusWebSocket) przenosi aktualizacje stanu zadań i dzienników do przeglądarki dzienników, natomiast Server-Sent Events (SseAlertController pod /api/v1/monitor/sse) przenoszą alerty i metryki. Gdy uruchomienie wytworzy nieprawidłowe dane, service/QuarantineService.kt i rekordy data_quarantine (Flyway engine V7) wstrzymują je od konsumpcji w dół strumienia, dopóki opiekun danych nie zatwierdzi, odrzuci lub nie edytuje ich ze strony DataQuarantine.


Etap 6 — Przechwytywanie pochodzenia danych

W trakcie wykonywania uruchomienia lineage/OpenLineageEmitter.kt z pipeline-engine emituje zdarzenia OpenLineage RunEvent. lineage-service pobiera je za pośrednictwem swojego OpenLineageEventController pod /api/v1/lineage/openlineage i buduje pochodzenie danych na poziomie zbiorów danych i kolumn.

  pipeline-engine                       lineage-service
  ───────────────                       ───────────────
  OpenLineageEmitter ──RunEvent──►  OpenLineageEventController

                                    LineageGraphBuilder

                              ┌───────────┼────────────┐
                              ▼           ▼            ▼
                       DatasetLineage  ColumnLineage  ImpactAnalyzer
                              │           │            │
                              └─────  PropagationWorker ┘
                                  (tags / PII propagation)

lineage-service ma wyłączone Flyway — ponownie wykorzystuje tabele pochodzenia danych, które tworzy metadata-service (Flyway metadata V4 i V48). Rejestruje pochodzenie danych na poziomie zbiorów danych i kolumn, wspiera zapytania podróży w czasie findLineageAsOf, uruchamia ImpactAnalyzer do analizy wpływu w dół strumienia oraz używa PropagationWorker do propagowania tagów i klasyfikacji PII w grafie. Graf pochodzenia danych jest prezentowany na stronie LineageExplorer.


Etap 7 — Zarządzanie i certyfikacja

Zbiory danych i potoki danych są zarządzane przez metadata-service. Zarządzanie to miejsce, gdzie dane są certyfikowane jako zdatne do użytku.

AktywnośćMechanizmUI
Przegląd i zatwierdzaniePolicyEngine, ApprovalWorkflow, przeglądy zarządzaniaReviewQueue
Kontrakty danych/api/v1/contracts, Flyway metadata V16DataContracts
Monitorowanie jakościQualityController, reguły / wyniki / oceny jakościQualityMonitoring
Rekomendacje/api/v1/governance/endorsementscentrum zarządzania
Słownik biznesowyencje słownika, Flyway metadata V10BusinessGlossary
Ewolucja schematu/api/v1/governance/schema-changesSchemaEvolution
Ślad audytuniezmienny dziennik audytu połączony łańcuchem skrótów, AuditChainVerifierServiceAuditTrail
ZgodnośćGdprService, PolishComplianceService, DSAR, maskowanieComplianceDashboard, DsarRequestsPage

Dziennik audytu jest niezmiennym zapisem połączonym łańcuchem skrótów (Flyway metadata V27 i V29) weryfikowanym przez AuditChainVerifierService. PII jest chronione przez cały czas: PiiMaskingFilter z bramy maskuje odpowiedzi, metadata-service uruchamia DynamicMaskingService oraz polityki maskowania (Flyway metadata V17), a klasyfikator PII zasila propagację PII w lineage-service.

Żądania dostępu osób, których dane dotyczą, w ramach RODO przebiegają według własnego pod-cyklu życia: zgłaszane jest DSAR (gdpr_dsar_requests, Flyway metadata V1), a usunięcie jest orkiestrowane przez erasure/ErasureOrchestrator.kt z pipeline-engine, który rejestruje erasure_runs i erasure_steps (Flyway engine V9).


Etap 8 — Wycofanie

Potok danych, który nie jest już potrzebny, zostaje wycofany. Jego definicja i historia wersji pozostają w metadata-service na potrzeby audytu, natomiast harmonogram jest wstrzymywany lub usuwany za pośrednictwem SchedulerController, a potok danych nie kwalifikuje się już do wykonania.

Wycofanie danych jest zarządzane oddzielnie. RetentionPolicyEnforcer w metadata-service stosuje polityki retencji, usunięcie RODO usuwa dane osobowe na żądanie, a rynek danych (Flyway metadata V20) może wycofać produkt danych, tak aby konsumenci już go nie odnajdywali. Historyczne rekordy uruchomień i pochodzenie danych są zachowywane do analizy wpływu i raportowania zgodności nawet po wycofaniu potoku danych.

  ACTIVE ──► DEPRECATED ──► RETIRED ──► (definition + history retained for audit)
              │                │
       schedule paused   removed from execution;
       consumers warned  data product withdrawn;
                         retention policy applies

Podsumowanie stanów potoku danych

Łącząc etapy razem, definicja potoku danych przechodzi przez te stany.

  DRAFT ──► VALIDATED ──► SCHEDULED ──► ACTIVE ──► DEPRECATED ──► RETIRED
   │            │            │           │
   │            │            │           └─ each trigger spawns a RUN
   │            │            │              (QUEUED → RUNNING → SUCCESS/FAILED)
   │            │            └─ cron / backfill / dependency / orchestrator
   │            └─ DSL + DAG + contract validation
   └─ authored in Design Studio / wizard / migration / template
Stan potoku danychZnaczenie
DRAFTUtworzony, jeszcze niezwalidowany; zapisany jako pipeline_version.
VALIDATEDPrzeszedł walidację DSL, DAG i kontraktów; uruchamialny na żądanie.
SCHEDULEDPowiązany z harmonogramem cron, uzupełnieniem lub łańcuchem zależności.
ACTIVEW regularnej eksploatacji; każde wyzwolenie wytwarza uruchomienie.
DEPRECATEDOznaczony do usunięcia; harmonogram wstrzymany, konsumenci ostrzeżeni.
RETIREDJuż niewykonywany; definicja i historia zachowane na potrzeby audytu.

Cykl życia migracji bazy danych Flyway

Schemat bazy danych platformy rozwija się poprzez migracje Flyway. Współdzielona baza danych PostgreSQL jest logicznie podzielona przez historie migracji poszczególnych usług.

UsługaLokalizacja FlywayTabela historiiZakres
metadata-servicedb/migration/metadataflyway_schema_historyV1–V54
pipeline-enginedb/migration/engineflyway_schema_history_engineV1–V9
monitor-servicedb/migration/monitorflyway_schema_history_monitorV9–V25
lineage-service— Flyway wyłączone —ponownie wykorzystuje tabele metadata

Migracja przebiega tym cyklem życia od utworzenia do aktywnego schematu:

  WRITE  ──►  COMMIT  ──►  STARTUP  ──►  APPLY  ──►  RECORD  ──►  VALIDATE
  V__NN     in db/         service       run         insert      ddl-auto
  numbered  migration/     boots         pending     row into    checks the
  SQL file  {service}/                   migrations  history     live schema
                                                     table       matches
  1. Write — programista dodaje numerowany plik SQL V__NN w folderze db/migration usługi.
  2. Commit — plik jest dostarczany wraz z kompilacją usługi.
  3. Startup — przy starcie usługi Flyway porównuje pliki migracji z tabelą historii usługi.
  4. Apply — oczekujące migracje uruchamiają się po kolei.
  5. Record — każda zastosowana migracja jest zapisywana jako wiersz w tabeli historii danej usługi.
  6. Validate — metadata-service uruchamia JPA ddl-auto: validate; pipeline-engine i lineage-service używają ddl-auto: none.

Permisywne ustawienia Flyway

Wdrożenie compose uruchamia Flyway z BASELINE_ON_MIGRATE, OUT_OF_ORDER, REPAIR_ON_MIGRATE, VALIDATE_ON_MIGRATE: false oraz BASELINE_VERSION: 20.1. Pozwalają one współdzielonej bazie danych z nakładającymi się, historycznie rozbieżnymi historiami migracji zbiec się przy starcie — ale stanowią ryzyko kruchości, a kolejność migracji ma znaczenie, ponieważ lineage-service zależy od tego, czy metadata-service zastosował już swoje migracje pochodzenia danych (V4 i V48).


Cykl życia wydania oprogramowania

Sama platforma jest dostarczana poprzez cykl kompilacja → test → wdrożenie → monitorowanie.

  BUILD ──────► TEST ──────► DEPLOY ──────► MONITOR
  Gradle JVM    vitest /     docker compose   Prometheus
  modules +     playwright   up on the VPS    scrape →
  Vite frontend / msw +      (git archive,    Grafana;
  + Python      backend      Flyway on        SSE alerts
  services      tests        startup)
FazaCo się dzieje
BuildUsługi JVM współdzielą jeden platform/Dockerfile z argumentem kompilacji BUILD_MODULE wybierającym moduł Gradle. Frontend Vite i dwie usługi Python mają własne konteksty kompilacji.
TestFrontend uruchamia testy jednostkowe vitest, testy end-to-end playwright (z obejściem VITE_AUTH_MODE=dev) oraz testy API z mockami msw; usługi backendu uruchamiają własne zestawy testów.
DeployDziałający system to pojedynczy VPS Debian aprowizowany przez deploy/deploy_to_vps.py, który łączy się przez SSH, rozpakowuje git archive i uruchamia docker compose up. Migracje Flyway uruchamiają się przy starcie usługi; nginx terminuje TLS.
MonitorPrometheus pobiera /actuator/prometheus oraz /metrics, zasilając pulpity Grafana; monitor-service emituje alerty przez SSE.

Dwie ścieżki wdrożenia

Repozytorium dokumentuje ścieżkę GitHub Actions → GCP Artifact Registry → GKE Autopilot z wdrożeniem typu canary, ale system faktycznie działający na produkcji to topologia Docker Compose na pojedynczym VPS. Cykl życia wydania opisany tutaj odzwierciedla działającą ścieżkę VPS.


Złożenie w całość

Trzy cykle życia zazębiają się: potok danych jest tworzony jako YAML, walidowany w przeglądarce i na serwerze, harmonogramowany przez cron lub uzupełnienie, wykonywany jako uporządkowane topologicznie uruchomienie DAG, monitorowany w sposób ciągły przez monitor-service, śledzony przez lineage-service jako zdarzenia OpenLineage, zarządzany i certyfikowany przez metadata-service, a na koniec wycofywany z zachowaną historią na potrzeby audytu. Pod spodem migracje Flyway rozwijają schemat przy każdym starcie usługi, a sama platforma przechodzi przez kompilację, test, wdrożenie i monitorowanie przy każdym wydaniu.

Poprzednia
Katalog konektorów