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 życia | Co przez niego przepływa | Zarządzany przez |
|---|---|---|
| Cykl życia potoku danych | Definicja potoku danych — od wersji roboczej do wycofania | metadata-service + pipeline-engine |
| Cykl życia uruchomienia | Pojedyncze wykonanie potoku danych | pipeline-engine, obserwowany przez monitor-service |
| Cykl życia danych | Zbiory danych wytwarzane przez potok — skatalogowane, śledzone, zarządzane | metadata-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ścia | Narzędzie | Uwagi |
|---|---|---|
| Projektowanie wizualne | Strona DesignStudio, kanwa React Flow | Węzły metodą przeciągnij i upuść; kanwa serializuje się do YAML |
| Tworzenie z przewodnikiem | CreatePipelineWizard (/pipelines/new) | Kreator krok po kroku dla typowych kształtów |
| Import migracji | migration-engine | Konwertuje starsze zadanie Informatica / Alteryx / SSIS / DataStage na DataFlow YAML |
| Instancjonowanie szablonu | Galeria PipelineTemplates | Rozpoczyna 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.
- Po stronie klienta. Usługa SPA
services/pipelineValidator.tssprawdza definicję w przeglądarce w trakcie edycji przez użytkownika, dając natychmiastową informację zwrotną w Studiu Projektowym. - Po stronie serwera. Gdy zażądane jest uruchomienie,
validation/PipelineValidator.ktz pipeline-engine waliduje sparsowaną definicję, avalidation/ParameterResolver.ktrozwią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.
| Komponent | Rola |
|---|---|
scheduler/CronScheduler.kt | Powtarzające się harmonogramy oparte na cron |
scheduler/BackfillManager.kt | Zadania uzupełniające na partycjach historycznych |
scheduler/DependencyChain.kt | Porządkowanie zależności między potokami danych |
scheduler/RetryPolicy.kt | Zachowanie 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 uruchomienia | Znaczenie |
|---|---|
QUEUED | Uruchomienie zostało zażądane i oczekuje na slot wykonawczy. |
RUNNING | DAG jest wykonywany; zadania działają w porządku topologicznym, równolegle w obrębie poziomu. |
SUCCESS | Wszystkie zadania zakończyły się pomyślnie. |
FAILED | Zadanie zawiodło i nie udało się go odzyskać. |
CANCELLED | Uruchomienie zostało anulowane kooperacyjnie przez ExecutionContext. |
HEALING | SelfHealingService 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.
| Zagadnienie | Komponent monitor-service | Uwidaczniane pod |
|---|---|---|
| Status i dzienniki uruchomienia | PipelineRunLogService, PipelineRunLogController | /monitor/runs, LogViewerPage |
| Alerty | AlertService, AlertDefinitionService, SseAlertController | /monitor/alerts |
| Koszt | PipelineCostService, CostAnomalyDetector | /monitor/costs |
| SLA / SLO | SlaBurnRateService | /monitor/sla |
| Świeżość | FreshnessTracker | /monitor/freshness |
| Wydajność | PerformanceMetricsService | /monitor/performance |
| Bazy odniesienia anomalii | AnomalyBaselineService | /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ść | Mechanizm | UI |
|---|---|---|
| Przegląd i zatwierdzanie | PolicyEngine, ApprovalWorkflow, przeglądy zarządzania | ReviewQueue |
| Kontrakty danych | /api/v1/contracts, Flyway metadata V16 | DataContracts |
| Monitorowanie jakości | QualityController, reguły / wyniki / oceny jakości | QualityMonitoring |
| Rekomendacje | /api/v1/governance/endorsements | centrum zarządzania |
| Słownik biznesowy | encje słownika, Flyway metadata V10 | BusinessGlossary |
| Ewolucja schematu | /api/v1/governance/schema-changes | SchemaEvolution |
| Ślad audytu | niezmienny dziennik audytu połączony łańcuchem skrótów, AuditChainVerifierService | AuditTrail |
| Zgodność | GdprService, PolishComplianceService, DSAR, maskowanie | ComplianceDashboard, 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 danych | Znaczenie |
|---|---|
DRAFT | Utworzony, jeszcze niezwalidowany; zapisany jako pipeline_version. |
VALIDATED | Przeszedł walidację DSL, DAG i kontraktów; uruchamialny na żądanie. |
SCHEDULED | Powiązany z harmonogramem cron, uzupełnieniem lub łańcuchem zależności. |
ACTIVE | W regularnej eksploatacji; każde wyzwolenie wytwarza uruchomienie. |
DEPRECATED | Oznaczony do usunięcia; harmonogram wstrzymany, konsumenci ostrzeżeni. |
RETIRED | Już 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ługa | Lokalizacja Flyway | Tabela historii | Zakres |
|---|---|---|---|
| metadata-service | db/migration/metadata | flyway_schema_history | V1–V54 |
| pipeline-engine | db/migration/engine | flyway_schema_history_engine | V1–V9 |
| monitor-service | db/migration/monitor | flyway_schema_history_monitor | V9–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
- Write — programista dodaje numerowany plik SQL
V__NNw folderzedb/migrationusługi. - Commit — plik jest dostarczany wraz z kompilacją usługi.
- Startup — przy starcie usługi Flyway porównuje pliki migracji z tabelą historii usługi.
- Apply — oczekujące migracje uruchamiają się po kolei.
- Record — każda zastosowana migracja jest zapisywana jako wiersz w tabeli historii danej usługi.
- 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)
| Faza | Co się dzieje |
|---|---|
| Build | Usł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. |
| Test | Frontend 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. |
| Deploy | Dział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. |
| Monitor | Prometheus 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.