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: copilot i migration-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ługaTechnologiaPort konteneraPort hostaDB / FlywayRola
api-gatewayKotlin / Spring Cloud Gateway (reaktywny WebFlux)80808085brak (używa Redis)Pojedynczy punkt wejścia: routing, walidacja JWT, ograniczanie szybkości, maskowanie PII, nagłówki bezpieczeństwa
metadata-serviceKotlin / Spring MVC + JPA80808181dataflow_metadata, Flyway metadata (V1–V54)Katalog, połączenia, nadzór, RODO, jakość, produkty danych, serwer MCP
pipeline-engineKotlin / Spring MVC + JPA80808082współdzielona DB, Flyway engine (V1–V9)Kompilacja/wykonywanie DAG potoków danych, harmonogramowanie, samonaprawa, usuwanie danych RODO
lineage-serviceKotlin / Spring MVC + JPA80808083współdzielona DB, Flyway wyłączonyPochodzenie danych na poziomie zbioru/kolumny, ingerencja OpenLineage, analiza wpływu
monitor-serviceKotlin / Spring MVC + JPA80808084współdzielona DB, Flyway monitor (V9–V25)Alerty, metryki, koszty, SLA, świeżość, powiadomienia, strumienie SSE
copilotPython 3.12 / FastAPI80008090współdzielona DB (pgvector)NL-to-pipeline, konwersacyjna AI, NL-to-SQL, wyszukiwanie RAG, RCA
migration-enginePython 3.12 / FastAPI80008091współdzielona DBKonwertuje starsze ETL (Informatica/Alteryx/SSIS/DataStage) → DataFlow YAML
frontendReact 19 / Vite, serwowany przez nginx803006nd.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

KomponentObrazPort(y)Przeznaczenie
PostgreSQLpgvector/pgvector:pg155432Współdzielony magazyn relacyjny + embeddingi pgvector; hostuje także bazę dataflow_keycloak
Keycloakquay.io/keycloak/keycloak:24.08180→8080Dostawca tożsamości OIDC / SSO; realm dataflow z realm-export.json
Kafkaconfluentinc/cp-kafka:7.7.19092 / 29092Kręgosłup strumieniowania zdarzeń (konektory, CDC/Debezium)
Zookeeperconfluentinc/cp-zookeeper:7.7.12181Koordynacja Kafki
Redisredis:7-alpine6379Ograniczanie szybkości bramy (przesuwne okno na sorted-set), buforowanie, magazyn sesji
MinIOminio/minio9000 / 9001Magazyn obiektowy zgodny z S3/GCS (konektory plików, CDR Parquet)
Prometheusprom/prometheus:v2.54.19090Pobieranie metryk
Grafanagrafana/grafana:11.2.23001→3000Pulpity 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

  1. RequestLoggingFilter — strukturalne logowanie dostępu.
  2. SecurityHeadersFilter — CSP i nagłówki bezpieczeństwa w odpowiedzi.
  3. PiiMaskingFilter — maskuje PII w ciałach odpowiedzi.
  4. RedisRateLimitFilter — per-trasowy limit przesuwnego okna Redis kluczowany rate_limit:{ip}:{endpoint}, z per-endpointowymi requestsPerMinute / burstCapacity; przy awarii Redis z fail-closed:true zwraca 503 + Retry-After.
  5. removeRequestHeader("Cookie") + preserveHostHeader().

Najważniejsze punkty routingu:

  • Standardowy REST → metadata / pipeline-engine / lineage / monitor.
  • Trasa WebSocket pipeline-run-websocket pod /api/v1/runs/*/stream przepisuje http→ws do 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-serviceNotificationInboxController znajduje 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ługaLokalizacja FlywayTabela historiiTryb DDL
metadata-servicedb/migration/metadata (V1–V54)flyway_schema_historyddl-auto: validate
pipeline-enginedb/migration/engine (V1–V9)flyway_schema_history_engineddl-auto: none
monitor-servicedb/migration/monitor (V9–V25)flyway_schema_history_monitor
lineage-serviceFlyway wyłączonyponownie 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) oraz StreamingQualityEngine / 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

WzorzecGdzie używanyMechanizm
Brama API / pojedynczy punkt wejściaCały ruch klientówSpring Cloud Gateway, deklaratywny RouteLocator
Uwierzytelnianie tokenowe + propagacja nagłówkówBrama → w dół stosuJWT Keycloak walidowany raz; nagłówki X-User-* wstrzyknięte
Serwer zasobów w głębokiej obronieKażda usługa JVMKażda uruchamia także serwer zasobów OAuth2 względem Keycloak
Synchroniczne HTTP między usługamiBrama ↔ usługiZwykłe HTTP po adresie URL z nazwą usługi
Strumieniowanie zdarzeńIngerencja konektorów, CDC, potoki strumienioweTematy Kafki; Debezium do przechwytywania zmian
Server-Sent Events (SSE)Alerty/metryki monitora, skrzynka powiadomień, czat copilotatext/event-stream; trasy SSE bramy; frontendowy SSEManager
WebSocketStatus uruchomienia potoku / strumieniowanie logów/api/v1/runs/*/stream, przepisanie http→ws bramy
Emisja/ingerencja OpenLineagepipeline-engine → lineage-serviceOpenLineageEmitter emituje RunEvent-y; OpenLineageEventController je przyjmuje
SDK konektorów / architektura wtyczekmetadata-service i pipeline-engineInterfejs ConnectorSDK; ConnectorRegistry poprzez ServiceLoader Javy
Adaptery zewnętrznych orkiestratorówpipeline-engineAirflowClient / AirflowDagGenerator, AutomateNowAdapter
Integracja z zewnętrznymi zasobami obliczeniowymipipeline-engineFlink (FlinkClusterManager) oraz Spark/Dataproc (DataprocJobSubmitter)
Abstrakcja podłączanych dostawców LLMcopilotABC LLMProvider + factory.py; LLM_PROVIDER przełącza anthropic / openrouter / local
RAG / wyszukiwanie wektorowecopilotEmbeddingi pgvector; degraduje płynnie, gdy niedostępne
Serwer Model Context Protocol (MCP)metadata-serviceMcpServerController pod /api/v1/mcp, /mcp
Pushdown / transpilacja SQLpushdown-sqlTranspiler oparty na JSqlParser, 8 dialektów docelowych
Kontrola wersji potoków oparta na Gitpipeline-engineGitRepositoryManager, PipelineVersionControl, YamlDiffEngine
Niezmienny log audytu z łańcuchem haszymetadata-serviceAuditChainVerifierService, łańcuch haszy logu audytu (Flyway V27/V29)

Zagadnienia przekrojowe

ZagadnienieImplementacja
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ściNagłówki X-User-* z bramy → ThreadLocal SecurityContextHolder
Ograniczanie szybkościPrzesuwne okno na sorted-set Redis per IP+endpoint; wyłącznik typu fail-closed; nagłówki X-RateLimit-*
Ochrona PIIBramowy PiiMaskingFilter na odpowiedziach; DynamicMaskingService w metadata-service; klasyfikator PII; wykrywanie propagacji PII w pochodzeniu danych
Audytcommon/AuditInterceptor + AuditPersistence; niezmienny log audytu z łańcuchem haszy w metadata-service z AuditChainVerifierService
Obserwowalność — metrykiSpring Actuator /actuator/prometheus; monitor-service OpenTelemetry + Prometheus; copilot udostępnia /metrics; pobierane przez Prometheus → Grafana
Obserwowalność — śledzeniemonitor-service OpenTelemetryConfig + TracingFilter; propagacja request-id (X-Request-Id)
Sprawdzenia kondycjiBramowy HealthController z actuator/health/{liveness,readiness,startup}; sondy gotowości badają .well-known Keycloak + kondycję metadata-service
Obsługa błędówcommon/GlobalExceptionHandler @ControllerAdvice + typowana hierarchia ApiException; SPA react-error-boundary + normalizacja błędów axios
Czas rzeczywistySSE (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-reads domyślnie ma wartość true w 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 to copilot / migration-engine; uzgadniane wyłącznie przez zmienne środowiskowe COPILOT_SERVICE_URL / MIGRATION_SERVICE_URL.
Poprzednia
Cykl życia danych i potoków