Architektura
Architektura systemu
DataFlow AI to produkt do integracji i modernizacji danych ETL/ELT zbudowany dla Polkomtela, polskiego operatora telekomunikacyjnego marki „Plus". Jest to wielojęzyczne monorepozytorium mikrousług osłonięte pojedynczą bramą API, a ta strona mapuje cały system od początku do końca.
Przegląd
DataFlow AI konwertuje, uruchamia, monitoruje i nadzoruje potoki danych dla telekomunikacyjnego zasobu danych. Platforma jest zorganizowana jako zestaw niezależnie budowanych usług, które współpracują ze sobą poprzez HTTP, Kafkę i współdzieloną bazę danych PostgreSQL. Istnieją cztery rodziny środowisk uruchomieniowych:
- Platforma JVM — Kotlin + Spring Boot, wielomodułowy build Gradle pod korzeniem pakietu
com.polkomtel.dataflow. Siedem modułów: pięć usług przeznaczonych do wdrożenia plus dwie biblioteki. - Usługi AI — dwie usługi Python 3.12 / FastAPI:
copilotimigration-engine. - Frontend — aplikacja jednostronicowa React 19 + Vite 7 + TypeScript.
- Narzędzia brzegowe — narzędzie wiersza poleceń Go (
backend/cli/) oraz rozszerzenie przeglądarki (browser-extension/).
Styl architektoniczny to mikrousługi za bramą API: pojedyncza reaktywna Spring Cloud Gateway jest jedynym punktem wejścia, uwierzytelnia każde żądanie względem Keycloak i proxuje do sześciu usług w dół stosu. Wywołania między usługami są synchroniczne po HTTP; Kafka przenosi obciążenia konektorów oraz przechwytywania zmian danych (CDC); SSE i WebSocket wypychają aktualizacje czasu rzeczywistego do interfejsu użytkownika.
Jedna baza danych, wiele historii migracji
Wszystkie usługi JVM i AI wskazują na tę samą instancję PostgreSQL i bazę danych (dataflow_metadata). Izolacja logiczna jest osiągana przez per-usługowe historie migracji Flyway oraz prefiksy nazw tabel (engine_*, monitor_*), a nie przez osobne fizyczne schematy.
Topologia mikrousług
Poniższy diagram pokazuje jednostki przeznaczone do wdrożenia, pojedynczy punkt wejścia oraz współdzieloną warstwę infrastruktury.
┌─────────────────────────────┐
│ Browser SPA (React 19) │
│ served by nginx :3006 │
└──────────────┬──────────────┘
│ HTTPS (Bearer JWT, axios baseURL=/api/v1)
│ + Keycloak OIDC login redirect
▼
┌─────────────────────────────┐
│ API Gateway :8085→8080 │
│ (Spring Cloud Gateway, │
│ reactive WebFlux) │
│ - JWT validation (Keycloak) │
│ - AuthFilter → X-User-* hdrs│
│ - Redis rate limiting │
│ - PII masking / sec headers │
└──┬───┬───┬───┬───┬───┬───────┘
/api/v1/... │ │ │ │ │ │
┌──────────────┬──────────┘ │ │ │ │ └──────────────┐
▼ ▼ ▼ ▼ ▼ ▼ ▼
┌─────────┐ ┌──────────┐ ┌──────────┐ ┌─────────┐ ┌──────────┐ ┌────────────┐
│metadata │ │ pipeline │ │ lineage │ │ monitor │ │ copilot │ │ migration │
│ :8181 │ │ engine │ │ :8083 │ │ :8084 │ │ :8090 │ │ engine │
│ │ │ :8082 │ │ │ │ │ │ (FastAPI)│ │ :8091 │
└────┬────┘ └────┬─────┘ └────┬─────┘ └────┬────┘ └────┬─────┘ └─────┬──────┘
│ │ │ │ │ │
└───────────┴────────────┴────────────┴───────────┴─────────────┘
│
┌────────────┬───────────┼────────────┬──────────────┐
▼ ▼ ▼ ▼ ▼
┌──────────┐ ┌──────────┐ ┌────────┐ ┌──────────┐ ┌──────────────┐
│PostgreSQL│ │ Kafka │ │ Redis │ │ MinIO │ │ Keycloak │
│ +pgvector│ │+Zookeeper│ │ │ │ (S3) │ │ (OIDC) │
└──────────┘ └──────────┘ └────────┘ └──────────┘ └──────────────┘
│
Prometheus → Grafana (scrape /actuator/prometheus, /metrics)
RouteConfig bramy odwołuje się do sześciu celów w dół stosu. Zwróć uwagę, że copilot-service i migration-service to nazwy tras bramy; kontenery Compose nazywają się copilot i migration-engine, a adresowane są przez zmienne środowiskowe COPILOT_SERVICE_URL / MIGRATION_SERVICE_URL.
Inwentarz usług i porty
Wewnątrz sieci Docker każda usługa JVM nasłuchuje na porcie 8080, a usługi AI na 8000; różnią się jedynie porty publikowane na hoście. Brama rozwiązuje usługi w dół stosu po adresach URL z nazwami usług.
| Usługa | Technologia | Port kontenera | Port hosta | DB / Flyway | Rola |
|---|---|---|---|---|---|
| api-gateway | Kotlin / Spring Cloud Gateway (reaktywny WebFlux) | 8080 | 8085 | brak (używa Redis) | Pojedynczy punkt wejścia: routing, walidacja JWT, ograniczanie szybkości, maskowanie PII, nagłówki bezpieczeństwa |
| metadata-service | Kotlin / Spring MVC + JPA | 8080 | 8181 | dataflow_metadata, Flyway metadata (V1–V54) | Katalog, połączenia, nadzór, RODO, jakość, produkty danych, serwer MCP |
| pipeline-engine | Kotlin / Spring MVC + JPA | 8080 | 8082 | współdzielona DB, Flyway engine (V1–V9) | Kompilacja/wykonywanie DAG potoków danych, harmonogramowanie, samonaprawa, usuwanie danych RODO |
| lineage-service | Kotlin / Spring MVC + JPA | 8080 | 8083 | współdzielona DB, Flyway wyłączony | Pochodzenie danych na poziomie zbioru/kolumny, ingerencja OpenLineage, analiza wpływu |
| monitor-service | Kotlin / Spring MVC + JPA | 8080 | 8084 | współdzielona DB, Flyway monitor (V9–V25) | Alerty, metryki, koszty, SLA, świeżość, powiadomienia, strumienie SSE |
| copilot | Python 3.12 / FastAPI | 8000 | 8090 | współdzielona DB (pgvector) | NL-to-pipeline, konwersacyjna AI, NL-to-SQL, wyszukiwanie RAG, RCA |
| migration-engine | Python 3.12 / FastAPI | 8000 | 8091 | współdzielona DB | Konwertuje starsze ETL (Informatica/Alteryx/SSIS/DataStage) → DataFlow YAML |
| frontend | React 19 / Vite, serwowany przez nginx | 80 | 3006 | nd. | Interfejs użytkownika aplikacji jednostronicowej |
Dwa moduły to biblioteki wbudowane w usługi JVM, a nie wdrażane jako kontenery: connector-sdk (framework konektorów plus 21 implementacji konektorów) oraz common (współdzielone bezpieczeństwo, modele, obsługa wyjątków). pushdown-sql to transpiler dialektów SQL Spring Boot używany jako wstrzykiwalna usługa/biblioteka; nie jest podpięty w pliku Compose jako samodzielny kontener.
Komponenty infrastruktury
| Komponent | Obraz | Port(y) | Przeznaczenie |
|---|---|---|---|
| PostgreSQL | pgvector/pgvector:pg15 | 5432 | Współdzielony magazyn relacyjny + embeddingi pgvector; hostuje także bazę dataflow_keycloak |
| Keycloak | quay.io/keycloak/keycloak:24.0 | 8180→8080 | Dostawca tożsamości OIDC / SSO; realm dataflow z realm-export.json |
| Kafka | confluentinc/cp-kafka:7.7.1 | 9092 / 29092 | Kręgosłup strumieniowania zdarzeń (konektory, CDC/Debezium) |
| Zookeeper | confluentinc/cp-zookeeper:7.7.1 | 2181 | Koordynacja Kafki |
| Redis | redis:7-alpine | 6379 | Ograniczanie szybkości bramy (przesuwne okno na sorted-set), buforowanie, magazyn sesji |
| MinIO | minio/minio | 9000 / 9001 | Magazyn obiektowy zgodny z S3/GCS (konektory plików, CDR Parquet) |
| Prometheus | prom/prometheus:v2.54.1 | 9090 | Pobieranie metryk |
| Grafana | grafana/grafana:11.2.2 | 3001→3000 | Pulpity obserwowalności |
Przepływ żądań
Cały ruch klientów wchodzi przez bramę pod /api/v1/**. RouteConfig bramy deklaruje ziarno RouteLocator, które mapuje prefiksy ścieżek na adresy URL usług w dół stosu, a każda trasa stosuje łańcuch filtrów.
Routing brzegowy
RequestLoggingFilter— strukturalne logowanie dostępu.SecurityHeadersFilter— CSP i nagłówki bezpieczeństwa w odpowiedzi.PiiMaskingFilter— maskuje PII w ciałach odpowiedzi.RedisRateLimitFilter— per-trasowy limit przesuwnego okna Redis kluczowanyrate_limit:{ip}:{endpoint}, z per-endpointowymirequestsPerMinute/burstCapacity; przy awarii Redis zfail-closed:truezwraca 503 +Retry-After.removeRequestHeader("Cookie")+preserveHostHeader().
Najważniejsze punkty routingu:
- Standardowy REST → metadata / pipeline-engine / lineage / monitor.
- Trasa WebSocket
pipeline-run-websocketpod/api/v1/runs/*/streamprzepisujehttp→wsdo pipeline-engine. - Trasy SSE dla strumieniowania czatu copilota oraz alertów/metryk monitora (
/api/v1/monitor/sse/**). /api/v1/notifications/**jest celowo kierowane do monitor-service —NotificationInboxControllerznajduje się tam, a nie w metadata.- Usługi AI montują swoje routery pod trzema prefiksami (
/api/v1/ai/**,/api/v1/copilot/**,/api/copilot/**), aby brama i frontend mogły się do nich odwoływać niezależnie od konfiguracji.
Przeglądarka → brama → usługa → baza danych
Browser ──▶ nginx :3006 load SPA
Browser ──▶ Keycloak :8180 OIDC login, redirect back with code
Browser keycloak-js exchanges code → access/refresh JWT
Browser ──▶ API Gateway :8085 GET/POST /api/v1/... (Authorization: Bearer JWT)
Gateway OAuth2 resource server validates JWT vs Keycloak JWKS
Gateway AuthFilter maps roles → DataFlowRole, injects X-User-* headers
Gateway RedisRateLimitFilter checks rate_limit:{ip}:{endpoint}
Gateway ──▶ downstream service :8080
Service re-validates JWT + reads X-User-* into SecurityContextHolder
Service RBAC permission check → query PostgreSQL → response
Gateway PiiMaskingFilter + SecurityHeadersFilter ──▶ Browser
Frontend używa pojedynczej instancji axios (frontend/src/api/client.ts) z baseURL=/api/v1. Interceptor żądań oczekuje na obietnicę waitForAuthReady, aby żadne żądanie nie zostało wystrzelone, zanim token Keycloak będzie dostępny, oraz wstrzykuje token bearer i token CSRF. Interceptor odpowiedzi normalizuje błędy i implementuje przepływ 401 → ciche odświeżenie tokena → ponowienie pierwotnego żądania.
Przepływ danych i trwałość
Współdzielona baza danych, partycjonowane migracje
Wszystkie usługi celują w tę samą instancję PostgreSQL i bazę danych. Izolacja logiczna wynika z per-usługowych historii migracji Flyway.
| Usługa | Lokalizacja Flyway | Tabela historii | Tryb DDL |
|---|---|---|---|
| metadata-service | db/migration/metadata (V1–V54) | flyway_schema_history | ddl-auto: validate |
| pipeline-engine | db/migration/engine (V1–V9) | flyway_schema_history_engine | ddl-auto: none |
| monitor-service | db/migration/monitor (V9–V25) | flyway_schema_history_monitor | — |
| lineage-service | Flyway wyłączony | ponownie wykorzystuje tabele pochodzenia danych metadata (V4/V48) | ddl-auto: none |
To wzorzec współdzielona baza danych, schemat-per-usługę-przez-prefiks-tabel: tabele engine_* dla silnika, tabele monitor_* dla monitora oraz tabele katalogu/nadzoru/RODO należące do metadata. lineage-service jest czystym konsumentem tabel pochodzenia danych metadata. Keycloak używa osobnej bazy dataflow_keycloak na tej samej instancji Postgres.
Liberalne ustawienia Flyway
Ustawienia Flyway w Compose są celowo liberalne — BASELINE_ON_MIGRATE, OUT_OF_ORDER, REPAIR_ON_MIGRATE, VALIDATE_ON_MIGRATE: false, BASELINE_VERSION: 20.1 — tak aby współdzielona baza danych z nakładającymi się starszymi historiami uzgodniła się przy starcie.
Magazyn obiektowy i wektory
- MinIO zapewnia magazyn zgodny z S3/GCS dla konektorów plików oraz konwersji CDR Parquet.
- pgvector stanowi podstawę magazynu embeddingów RAG copilota. Copilot inicjalizuje pulę asyncpg i zapewnia schemat embeddingów przy starcie, degradując do
rag_mode: "unavailable", jeśli pgvector jest nieosiągalny.
Strumieniowanie zdarzeń
Kafka + Zookeeper tworzą kręgosłup strumieniowania. Wykorzystanie Kafki jest skoncentrowane w:
- connector-sdk — konektor Kafka (
KafkaConnector,KafkaOffsetManager,KafkaSchemaRegistry, SerDe), konektor Pub/Sub oraz CDC poprzez Debezium (DebeziumCdcConnector,DebeziumCdcManager,DebeziumEventConsumer). - pipeline-engine — zadania strumieniowe (
FlinkJobBuilder) orazStreamingQualityEngine/StreamingQualityController.
Usługi JVM otrzymują KAFKA_BOOTSTRAP_SERVERS: kafka:29092 (wewnętrzny listener). Kafka przenosi obciążenia ingerencji konektorów i CDC, zamiast działać jako szyna poleceń między mikrousługami — wywołania między usługami pozostają synchronicznym HTTP.
Wzorce integracji
| Wzorzec | Gdzie używany | Mechanizm |
|---|---|---|
| Brama API / pojedynczy punkt wejścia | Cały ruch klientów | Spring Cloud Gateway, deklaratywny RouteLocator |
| Uwierzytelnianie tokenowe + propagacja nagłówków | Brama → w dół stosu | JWT Keycloak walidowany raz; nagłówki X-User-* wstrzyknięte |
| Serwer zasobów w głębokiej obronie | Każda usługa JVM | Każda uruchamia także serwer zasobów OAuth2 względem Keycloak |
| Synchroniczne HTTP między usługami | Brama ↔ usługi | Zwykłe HTTP po adresie URL z nazwą usługi |
| Strumieniowanie zdarzeń | Ingerencja konektorów, CDC, potoki strumieniowe | Tematy Kafki; Debezium do przechwytywania zmian |
| Server-Sent Events (SSE) | Alerty/metryki monitora, skrzynka powiadomień, czat copilota | text/event-stream; trasy SSE bramy; frontendowy SSEManager |
| WebSocket | Status uruchomienia potoku / strumieniowanie logów | /api/v1/runs/*/stream, przepisanie http→ws bramy |
| Emisja/ingerencja OpenLineage | pipeline-engine → lineage-service | OpenLineageEmitter emituje RunEvent-y; OpenLineageEventController je przyjmuje |
| SDK konektorów / architektura wtyczek | metadata-service i pipeline-engine | Interfejs ConnectorSDK; ConnectorRegistry poprzez ServiceLoader Javy |
| Adaptery zewnętrznych orkiestratorów | pipeline-engine | AirflowClient / AirflowDagGenerator, AutomateNowAdapter |
| Integracja z zewnętrznymi zasobami obliczeniowymi | pipeline-engine | Flink (FlinkClusterManager) oraz Spark/Dataproc (DataprocJobSubmitter) |
| Abstrakcja podłączanych dostawców LLM | copilot | ABC LLMProvider + factory.py; LLM_PROVIDER przełącza anthropic / openrouter / local |
| RAG / wyszukiwanie wektorowe | copilot | Embeddingi pgvector; degraduje płynnie, gdy niedostępne |
| Serwer Model Context Protocol (MCP) | metadata-service | McpServerController pod /api/v1/mcp, /mcp |
| Pushdown / transpilacja SQL | pushdown-sql | Transpiler oparty na JSqlParser, 8 dialektów docelowych |
| Kontrola wersji potoków oparta na Git | pipeline-engine | GitRepositoryManager, PipelineVersionControl, YamlDiffEngine |
| Niezmienny log audytu z łańcuchem haszy | metadata-service | AuditChainVerifierService, łańcuch haszy logu audytu (Flyway V27/V29) |
Zagadnienia przekrojowe
| Zagadnienie | Implementacja |
|---|---|
| Uwierzytelnianie (AuthN) | Keycloak 24 OIDC, realm dataflow; keycloak-js w SPA; JWT walidowane na bramie i ponownie na każdym serwerze zasobów |
| Autoryzacja (AuthZ) | Hierarchiczne RBAC w common/RBACService: DataFlowRole ADMIN(100) > ENGINEER(75) > ANALYST(50) > STEWARD(40) > VIEWER(25); reguły metody/roli bramy (DELETE→ADMIN, zapis→ADMIN/ENGINEER, GET→dowolna rola) |
| Propagacja tożsamości | Nagłówki X-User-* z bramy → ThreadLocal SecurityContextHolder |
| Ograniczanie szybkości | Przesuwne okno na sorted-set Redis per IP+endpoint; wyłącznik typu fail-closed; nagłówki X-RateLimit-* |
| Ochrona PII | Bramowy PiiMaskingFilter na odpowiedziach; DynamicMaskingService w metadata-service; klasyfikator PII; wykrywanie propagacji PII w pochodzeniu danych |
| Audyt | common/AuditInterceptor + AuditPersistence; niezmienny log audytu z łańcuchem haszy w metadata-service z AuditChainVerifierService |
| Obserwowalność — metryki | Spring Actuator /actuator/prometheus; monitor-service OpenTelemetry + Prometheus; copilot udostępnia /metrics; pobierane przez Prometheus → Grafana |
| Obserwowalność — śledzenie | monitor-service OpenTelemetryConfig + TracingFilter; propagacja request-id (X-Request-Id) |
| Sprawdzenia kondycji | Bramowy HealthController z actuator/health/{liveness,readiness,startup}; sondy gotowości badają .well-known Keycloak + kondycję metadata-service |
| Obsługa błędów | common/GlobalExceptionHandler @ControllerAdvice + typowana hierarchia ApiException; SPA react-error-boundary + normalizacja błędów axios |
| Czas rzeczywisty | SSE (alerty/metryki/powiadomienia/czat copilota) + WebSocket (strumieniowanie uruchomień potoków) |
| Zgodność | Moduły RODO + zgodności z polskim prawem (GdprService, PolishComplianceService, DSAR, retencja, powiadamianie o naruszeniach); harmonogramy domyślnie ustawione na Europe/Warsaw |
Kluczowe przepływy sekwencji
Logowanie i uwierzytelnione żądanie API
Browser ──(1) load SPA──▶ nginx :3006
Browser ──(2) OIDC redirect──▶ Keycloak :8180
Keycloak ──(3) login + redirect back with code──▶ Browser
Browser (4) keycloak-js exchanges code → access/refresh JWT; waitForAuthReady resolves
Browser ──(5) GET /api/v1/... (Bearer JWT)──▶ API Gateway :8085
Gateway (6) OAuth2 resource server validates JWT signature/issuer/audience vs JWKS
Gateway (7) AuthFilter maps Keycloak roles → DataFlowRole, injects X-User-* headers
Gateway (8) RedisRateLimitFilter checks/updates rate_limit:{ip}:{endpoint}
Gateway ──(9) proxy to downstream :8080──▶ metadata / pipeline / monitor / ...
Service (10) re-validates JWT + reads X-User-* into SecurityContextHolder
Service (11) RBAC permission check → query PostgreSQL → response
Gateway (12) PiiMaskingFilter + SecurityHeadersFilter ──▶ Browser
On 401: response interceptor silently refreshes the token and retries the request.
Uruchomienie potoku i strumieniowanie statusu w czasie rzeczywistym
Frontend ──POST /api/v1/pipelines/{id}/run──▶ Gateway ──▶ pipeline-engine (ExecutionController)
pipeline-engine: PipelineRunner parses YAML DSL → ParameterResolver → PipelineValidator
→ DagBuilder builds ExecutionDAG (cycle/dangling detection)
→ executes tasks in topological order (parallel within DAG levels)
→ ExecutionContext supports cooperative cancellation
For streaming/compute: FlinkJobSubmitter / Spark DataprocJobSubmitter; or AirflowClient
On task failure: SelfHealingService → FailureClassifier → RecoveryStrategies
pipeline-engine: ExecutionEventPublisher / PipelineRunLogPublisher push run events
pipeline-engine: OpenLineageEmitter emits RunEvents ──▶ lineage-service /openlineage/events
Frontend ──WS /api/v1/runs/{runId}/stream──▶ Gateway (http→ws) ──▶ PipelineStatusWebSocket
→ live task-state / log updates streamed to the SPA log viewer
On completion: monitor-service ingests run metrics → cost/SLA/freshness eval → alerts
Żądanie do copilota (NL-to-pipeline / czat z RAG)
Frontend ──POST /api/v1/copilot/chat (SSE)──▶ Gateway ──▶ copilot :8090
copilot: RAGService queries pgvector for relevant catalog/pipeline context
(rag_max_results=5, similarity_threshold=0.72)
degrades to non-RAG chat if pgvector unavailable
copilot: LLM provider factory selects anthropic | openrouter | local per LLM_PROVIDER
copilot: LLMProvider.stream() → token deltas (StreamChunk) streamed back as SSE events
(X-Copilot-Confidence header set); can produce DataFlow YAML pipeline definitions
Frontend: AICopilotSidebar renders streamed tokens; generated pipeline → Design Studio
Ingerencja konektorów i pochodzenie danych
metadata-service / pipeline-engine load a connector via ConnectorRegistry (ServiceLoader)
Connector.discoverSchema() → catalog tables/columns persisted (metadata-service catalog)
Pipeline execution: Connector.extractData() → DataStream → transform → loadData() to target
Streaming/CDC connectors (Kafka, Debezium) publish/consume change events via Kafka topics
pipeline-engine OpenLineageEmitter → lineage-service builds dataset/column lineage graph
lineage-service: ImpactAnalyzer + PropagationWorker propagate tags/PII across the graph
Obserwacje architektoniczne
Kilka właściwości systemu warto mieć na uwadze przy jego rozszerzaniu:
- Współdzielona pojedyncza baza danych dla wszystkich usług JVM i Python. Izolacja odbywa się przez prefiksy tabel i per-usługowe tabele historii Flyway, a nie fizyczne schematy. Sprzęga to usługi na warstwie danych — bliżej rozproszonego monolitu nad jedną bazą danych niż w pełni autonomicznych mikrousług.
- lineage-service ma wyłączony Flyway i zależy od tego, czy metadata-service już zastosował swoje migracje pochodzenia danych (V4/V48) — ukryta zależność kolejności w czasie wdrażania.
- Komunikacja między usługami jest synchronicznym HTTP bez siatki usług, wyłączników czy odkrywania usług poza Docker DNS i adresami URL ze zmiennych środowiskowych.
- Tożsamość jest podwójnie walidowana na bramie i ponownie na każdym serwerze zasobów — silna głęboka obrona — ale flaga
dev-permit-readsdomyślnie ma wartośćtruew Compose, co zezwoliłoby na nieuwierzytelnione odczyty, gdyby pozostała włączona. - Niezgodność nazewnictwa usług: trasy bramy nazywają
copilot-service/migration-service, podczas gdy kontenery Compose tocopilot/migration-engine; uzgadniane wyłącznie przez zmienne środowiskoweCOPILOT_SERVICE_URL/MIGRATION_SERVICE_URL.