Architektura
Usługi platformowe backendu
Warstwa platformowa DataFlow AI to wielomodułowy build Gradle usług Kotlin + Spring Boot pod korzeniem pakietu com.polkomtel.dataflow. Ta strona dokumentuje każdy moduł — pięć usług przeznaczonych do wdrożenia oraz trzy moduły biblioteczne/narzędziowe.
Przegląd modułów
Każda usługa posiada własną tabelę historii Flyway (flyway_schema_history, ..._engine, ..._monitor), dzięki czemu współdzielona baza danych może hostować wiele zestawów migracji bez kolizji.
| Moduł | Port | Typ | DB / Flyway | Przeznaczenie |
|---|---|---|---|---|
api-gateway | 8080 | Spring Cloud Gateway (reaktywny WebFlux) | brak (Redis) | Routing, walidacja JWT, ograniczanie szybkości, maskowanie PII |
metadata-service | 8081 | Spring MVC (JPA) | własna DB, Flyway metadata (V1–V54) | Katalog, nadzór, RODO, jakość, metadane potoków |
pipeline-engine | 8082 | Spring MVC (JPA) | własna DB, Flyway engine (V1–V9) | Wykonywanie potoków, harmonogramowanie, samonaprawa |
pushdown-sql | 8083 | Spring Boot (bez DB/bezpieczeństwa) | brak | Biblioteka/usługa transpilacji dialektów SQL |
lineage-service | 8084 | Spring MVC (JPA) | odczytuje DB metadata, Flyway wyłączony | Śledzenie pochodzenia danych i graf |
monitor-service | 8080 (SERVER_PORT) | Spring MVC (JPA) | własna DB, Flyway monitor (V9–V25) | Monitoring, alerty, koszty, SLA, powiadomienia |
connector-sdk | nd. | Moduł biblioteczny | brak | Framework konektorów + 21 implementacji konektorów |
common | nd. | Moduł biblioteczny | brak | Współdzielone bezpieczeństwo, modele, obsługa wyjątków |
Usługi
api-gateway
Brama jest pojedynczym punktem wejścia do platformy. To reaktywna aplikacja Spring Cloud Gateway (GatewayApplication.kt, reaktywny typ web), która waliduje JWT Keycloak, propaguje tożsamość do usług w dół stosu poprzez nagłówki oraz stosuje ograniczanie szybkości, nagłówki bezpieczeństwa, maskowanie PII i logowanie żądań.
Port: 8080 (host 8085).
Kluczowe klasy
| Klasa | Rola |
|---|---|
config/RouteConfig.kt | Deklaratywne ziarno RouteLocator kierujące ścieżki /api/v1/** do sześciu usług w dół stosu. Każda trasa stosuje logowanie żądań, nagłówki bezpieczeństwa, maskowanie PII, limit szybkości Redis, removeRequestHeader("Cookie"), preserveHostHeader(). Zawiera trasę WebSocket pipeline-run-websocket (przepisanie http→ws) oraz trasy SSE dla czatu copilota i alertów/metryk monitora. |
config/SecurityConfig.kt | @EnableWebFluxSecurity. CSRF/basic/formLogin wyłączone; serwer zasobów OAuth2 z JWT. Reguły metody/roli: DELETE→ROLE_ADMIN; POST/PUT/PATCH→ROLE_ADMIN/ROLE_ENGINEER; GET→dowolna rola. Flaga dev-permit-reads zezwala na nieuwierzytelnione żądania GET oraz POST do copilota/wyszukiwania w trybie deweloperskim. |
config/ReactiveKeycloakJwtConverter.kt | Konwertuje roszczenia JWT Keycloak na uprawnienia Spring. |
filter/AuthFilter.kt | GlobalFilter (kolejność -90). Wyodrębnia roszczenia JWT, mapuje role Keycloak na DataFlowRole poprzez RBACService, wstrzykuje X-User-Id, X-User-Email, X-User-Display-Name, X-User-Role, X-Workspace-Id, X-User-Groups. |
filter/RedisRateLimitFilter.kt | Limiter przesuwnego okna na sorted-set Redis kluczowany rate_limit:{ip}:{endpoint}. Przy awarii Redis z fail-closed:true zwraca 503 + Retry-After; emituje nagłówki X-RateLimit-*. |
filter/PiiMaskingFilter.kt, SecurityHeadersFilter.kt, RequestLoggingFilter.kt, CorsConfig.kt | Maskowanie PII w ciele odpowiedzi, nagłówki bezpieczeństwa, strukturalne logowanie dostępu, CORS. |
controller/HealthController.kt | /actuator/health/{liveness,readiness,startup} oraz /api/v1/health. Sondy gotowości badają .well-known/openid-configuration Keycloak oraz /actuator/health metadata-service. |
Powierzchnia REST: brama proxuje ruch; jej jedynymi własnymi endpointami są trasy actuator/health HealthController. Wydawca/JWK OAuth2 pochodzą z KEYCLOAK_ISSUER_URI / KEYCLOAK_JWK_SET_URI; adresy URL usług w dół stosu oraz źródła CORS są sterowane zmiennymi środowiskowymi.
common
common to współdzielona biblioteka używana przez każdą usługę JVM. Dostarcza prymitywy bezpieczeństwa, modele domenowe oraz obsługę wyjątków — bez własnej aplikacji Spring.
Kluczowe klasy
security/RBACService.kt— hierarchiczne RBAC. EnumDataFlowRole: ADMIN(100) > ENGINEER(75) > ANALYST(50) > STEWARD(40) > VIEWER(25). EnumPermissionzrequiredRole(np.CREATE_PIPELINE→ENGINEER,DELETE_PIPELINE→ADMIN,USE_COPILOT→ANALYST).mapKeycloakRolesToDataFlowRole()mapuje grupy AD (PLK-BI-Admins), nazwy pisane małymi literami,dataflow-*oraz natywne dla realmu role Keycloak (org_admin,data_engineer,business_analyst). Pomocnicze:hasPermission,canAccessWorkspace,canModifyPipeline,canDeletePipeline,canRunPipeline.security/JwtAuthenticator.kt— dekoduje JWT, wymusza jawną walidację odbiorcyaud(domyślniedataflow-api), zwracaAuthenticatedPrincipal.security/KeycloakJwtConverter.kt— wyodrębnia nazwę principal + uprawnienia z JWT (wariant servletowy).security/SecurityContext.kt— klasa danychSecurityContext(userId, email, role, workspaceId, groups) zisAdmin/isEngineerOrAbove/isAnalystOrAbove; ThreadLocalSecurityContextHolderwypełniany z nagłówków bramy w usługach w dół stosu.security/SecurityConfig.kt,SecurityInterceptorConfig.kt— konfiguracja bezpieczeństwa servletu i podpięcie interceptorów.security/AuditInterceptor.kt,AuditPersistence.kt— przechwytywanie i trwałość zdarzeń audytu na poziomie żądania.exception/ApiException.kt,GlobalExceptionHandler.kt— typowane wyjątki (NotFoundException,ForbiddenException) plus globalny handler@ControllerAdvice.model/— współdzielone modele domenowe:Pipeline,PipelineRun,Connection,User,Workspace,Alert,AuditEvent,enums.kt.config/JacksonConfig.kt— konfiguracja Jackson.
metadata-service
metadata-service to największy moduł — centralny kręgosłup katalogu i nadzoru. Zarządza połączeniami, metadanymi potoków, katalogiem danych, glosariuszem, rejestrem schematów, klasyfikacją, przepływami nadzoru, RODO i zgodnością z polskim prawem, jakością danych, produktami danych i marketplace, tagami, serwerem MCP oraz rejestrem AI. Jest oznaczony adnotacjami @EnableScheduling + @EnableCaching i skanuje podpakiety, w tym catalog, governance, gdpr, search, policy, quality, domain, compliance, aiRegistry, workflow, tags oraz mcp.
Port: 8081. JPA ddl-auto: validate; pula Hikari maks. 20 / 5 bezczynnych; Jackson SNAKE_CASE.
Powierzchnia REST (kontrolery według ścieżki bazowej)
| Obszar | Ścieżka(-i) bazowa |
|---|---|
| Połączenia / konektory | /api/v1/connections, /api/v1/connectors, /api/v1/marketplace |
| Potoki / szablony | /api/v1/pipelines, /api/v1/templates, /api/v1/runbooks |
| Katalog / jakość / profilowanie | /api/v1/catalog, /api/v1/quality, /api/v1/cde (krytyczne elementy danych) |
| Nadzór | /api/v1/governance oraz /access-reviews, /contracts, /endorsements, /reviews, /schema-changes, /tags |
| RODO / zgodność | /api/v1/gdpr, /api/v1/dsar, /api/v1/compliance, /api/v1/compliance/dpia, /api/v1/admin/abac/policy |
| Produkty danych / wersje / domeny | /api/v1/data-products, /api/v1/data-versions, /api/v1/domains, /api/v1/contracts |
| Audyt | /api/v1/audit-log, /api/v1/admin/audit |
| Maskowanie | /api/v1/masking |
| Incydenty / RCA | /api/v1/incidents |
| CDR | /api/v1/cdr (telekomunikacyjne metryki Call Detail Record) |
| Klasyfikacja | /api/v1/classification/rules |
| Przepływy pracy | /api/v1/workflows |
| MCP | /api/v1/mcp, /mcp |
| Rejestr AI | /api/v1/admin/ai-registry |
| Wyszukiwanie | /api/v1/search (globalne wyszukiwanie po widoku zmaterializowanym) |
| Administracja | /api/v1/admin/users, /api/v1/admin/workspaces, /api/v1/admin/settings |
| Przesyłanie plików | /api/v1/data |
Wśród godnych uwagi usług znajdują się PolicyEngine, ApprovalWorkflow, AccessReviewCycleService, GdprService, PolishComplianceService, RetentionPolicyEnforcer, DynamicMaskingService, IncidentRcaService, DataClassificationService, WorkflowEngine oraz AuditChainVerifierService (log audytu z łańcuchem haszy).
Model danych
Około 80 encji JPA znajduje się pod metadata/entity/ i jego podpakietami. Rdzeń (V1): workspaces, users, connections, pipelines, pipeline_versions, pipeline_runs, pipeline_tasks. Katalog (V9): catalog_databases/schemas/tables/columns/tags, relacje, użycia. RODO (V1): gdpr_dsar_requests, gdpr_consent_records, gdpr_erasure_reports, gdpr_audit_trail. Encje pochodzenia danych ColumnLineageEntity, DatasetLineageEntity, OpenLineageEventEntity oraz ManualLineageEdgeEntity należą do tego modułu i są współdzielone z lineage-service.
Migracje Flyway (db/migration/metadata, V1–V54)
| Wersja | Zawartość |
|---|---|
| V1–V5 | Rdzeń schematu, potoki, monitoring, pochodzenie danych, audyt |
| V9–V14 | Katalog, glosariusz, rejestr schematów, klasyfikacja, nadzór, przepływ pracy |
| V16–V21 | Kontrakty danych, dynamiczne maskowanie, kontrola wersji danych, marketplace danych, zgodność z polskim prawem (V20_1), trwałość audytu RODO |
| V27 / V29 | Łańcuch haszy logu audytu |
| V33 / V35 | Katalog marketplace konektorów, widok zmaterializowany wyszukiwania globalnego |
| V44–V54 | Ramy regulacyjne, RCA incydentów (V47), OpenLineage + ręczne krawędzie (V48), silnik przepływów pracy (V52), rejestr AI (V53), propagacja tagów + uruchomienia usuwania (V54) |
Ziarno deweloperskie znajduje się w db/dev/V6__seed_data.sql, a przykładowy szablon potoku w resources/pipeline-templates/subscriber-360.yaml.
pipeline-engine
pipeline-engine wykonuje potoki zdefiniowane w DSL YAML, harmonogramuje uruchomienia, zarządza wykonywaniem DAG, integruje zewnętrzne orkiestratory (Airflow, AutomateNow) i silniki obliczeniowe (Flink, Spark/Dataproc), wykonuje samonaprawę, orkiestruje usuwanie danych RODO, poddaje kwarantannie błędne dane oraz zapewnia kontrolę wersji potoków opartą na Git. Klasa aplikacji: PipelineEngineApplication.kt.
Port: 8082. JPA ddl-auto: none; Flyway db/migration/engine, historia flyway_schema_history_engine; strefa czasowa harmonogramu domyślnie Europe/Warsaw.
Kluczowe klasy
- Wykonywanie —
execution/PipelineRunner.ktorkiestruje uruchomienie: parsuj YAML → rozwiąż parametry → waliduj → zbuduj DAG → wykonaj zadania w kolejności topologicznej (równolegle w obrębie poziomów DAG, pula wątków stałej wielkości) → agreguj wyniki, ze współpracującym anulowaniem poprzezExecutionContext. Wspierany przezTaskExecutor.ktorazExecutionContext.kt. - DAG —
dag/DagBuilder.kt,dag/ExecutionDAG.kt,dag/DagNode.kt, z wykrywaniem cykli i zwisających zależności. - DSL —
dsl/PipelineYamlParser.kt,dsl/PipelineYamlSchema.kt. - Walidacja —
validation/PipelineValidator.kt,validation/ParameterResolver.kt(zmienne wbudowane, zmienne środowiskowe, parametry uruchomieniowe). - Warstwa usług —
service/PipelineCompiler.kt,ExecutionPlanner.kt,TaskScheduler.kt,PipelineRunLogPublisher.kt. - Harmonogramowanie —
scheduler/CronScheduler.kt,BackfillManager.kt,DependencyChain.kt,RetryPolicy.kt. - Orkiestratory —
orchestrator/AirflowClient.kt,AirflowDagGenerator.kt,AutomateNowAdapter.kt,WebhookCallbackService.kt. - Silniki obliczeniowe —
flink/(FlinkClusterManager,FlinkJobBuilder/Submitter,FlinkCheckpointManager,FlinkSqlBridge) orazspark/(SparkClusterManager,SparkJobBuilder,DataprocJobSubmitter,SparkSqlBridge). - Samonaprawa —
healing/SelfHealingService.kt,FailureClassifier.kt,RecoveryStrategies.kt. - Git —
git/GitRepositoryManager.kt,PipelineVersionControl.kt,YamlDiffEngine.kt,GitWebhookHandler.kt. - Usuwanie (RODO) —
erasure/ErasureOrchestrator.kt,ErasureController. - Jakość —
quality/DataQuarantineService.kt,streaming/StreamingQualityEngine.kt. - Czas rzeczywisty —
realtime/PipelineStatusWebSocket.kt,ExecutionEventPublisher.kt. - Pochodzenie danych —
lineage/OpenLineageEmitter.kt.
Powierzchnia REST
| Kontroler | Ścieżka bazowa | Kluczowe operacje |
|---|---|---|
ExecutionController | /api/v1 | POST /pipelines/{id}/run, GET /runs/{runId}, POST /runs/{runId}/cancel, POST /runs/{runId}/status |
SchedulerController | /api/v1/scheduler | CRUD harmonogramów + pause/resume, /next-runs, CRUD backfill, CRUD zależności |
OrchestratorController | /api/v1/orchestrator | /trigger, /callback, dagi/sync/trigger/status Airflow, webhooki |
GitController | /api/v1/git | repozytoria, historia, diff, rollback, gałęzie, merge, webhooki |
HealingController | /api/v1/pipelines | /{id}/healing-history, /healing-summary, /healing-policy |
QuarantineController | /api/v1/quarantine | list/get, approve/reject/edit-approve, operacje masowe, reguły auto-zwolnienia |
ErasureController | /api/v1/dsar | POST /{dsarId}/erasure/orchestrate, status uruchomienia usuwania |
StreamingQualityController | /api/v1/monitor/streaming | jakość per-potok, reguły, reset, SQL |
PipelineReportController | /api/v1/pipelines/{id}/report | raporty uruchomień potoków |
Model danych i migracje (db/migration/engine, V1–V9)
Encje są prefiksowane engine_*: engine_schedules, engine_backfill_jobs/partitions, engine_execution_runs (z kolumną JSONB task_states), engine_retry_states, engine_circuit_breakers, engine_dead_letter_entries, engine_running_pipelines, engine_webhook_registrations, engine_healing_attempts/policies. Bez prefiksu: data_quarantine, quarantine_auto_release_rules, erasure_runs, erasure_steps. Migracje: V1 harmonogramowanie, V3 łańcuch zależności, V4 uruchomienia wykonań, V5 webhook + audyt, V6 ponowienie/alert, V7 kwarantanna, V8 próby naprawy, V9 orkiestracja usuwania.
lineage-service
lineage-service śledzi pochodzenie danych — pochodzenie na poziomie zbioru i kolumny, analizę wpływu, ingerencję zdarzeń OpenLineage, budowę i wyszukiwanie grafu pochodzenia, pochodzenie narzędzi BI (Tableau/Looker/PowerBI), ręczne tworzenie pochodzenia oraz propagację tagów. Konsumuje RunEvent-y OpenLineage emitowane przez pipeline-engine. Klasa aplikacji: LineageServiceApplication.kt.
Port: 8084. JPA ddl-auto: none; flyway.enabled: false — polega na schemacie metadata-service (tabele pochodzenia danych tworzone są przez migracje metadata V4/V48).
Zależność kolejności w czasie wdrażania
Ponieważ lineage-service ma wyłączony Flyway, nie może startować, dopóki metadata-service nie zastosuje swoich migracji pochodzenia danych. Ta kolejność jest niejawna i nie jest wymuszana przez sam schemat.
Kluczowe klasy
service/LineageService.kt— główna usługa:getDatasetLineage,getColumnLineage,recordEvent/recordEvents,getUpstream/getDownstream,getSubgraph,getFullLineageGraph,getStats. Używa najpierwLineageGraphBuilder, wracając do magazynu w pamięciConcurrentHashMapplusLineagePersistence.service/LineageGraphBuilder.kt— buduje graf pochodzenia danych ze zdarzeń.service/ImpactAnalyzer.kt— analiza wpływu w dół stosu.service/ColumnLineageExtractor.kt— wyodrębnia pochodzenie na poziomie kolumny.service/LineageSearchService.kt— wyszukiwanie zbiorów/kolumn, gorące zbiory danych, wykrywanie propagacji PII.service/connectors/—BiLineageConnector,TableauLineageConnector,LookerLineageConnector,PowerBiLineageConnector.propagation/PropagationWorker.kt,TagPropagationLogEntity/Repository.kt— propagacja tagów wzdłuż pochodzenia danych.repository/JpaLineageRepositories.kt—DatasetLineageRepository,ColumnLineageRepository, ponownie wykorzystujące encje metadata-service; obsługujefindLineageAsOf(podróż w czasie) oraz filtrowanie PII.model/LineageEvent.kt—RunEventw stylu OpenLineage,LineageGraph,LineageNode,LineageJobNode,SchemaField.
Powierzchnia REST
| Kontroler | Ścieżka bazowa | Kluczowe operacje |
|---|---|---|
LineageController | /api/v1/lineage | GET /datasets/{id}, GET /columns/{id}, GET /impact/{id}, POST /events, GET /graph/{datasetId}, GET /upstream/{datasetId}, GET /downstream/{datasetId}, GET /search, GET /stats, GET /subgraph/{datasetId}, GET /hot-datasets, GET /pii/{datasetId}/{columnName} |
OpenLineageEventController | /api/v1/lineage/openlineage | POST /events, GET /events |
BiLineageController | /api/v1/lineage/bi | POST /ingest, POST /{tool}/ingest |
LineageAuthoringController | /api/v1/lineage/manual-edges | GET, POST, PATCH /{id}, DELETE /{id} |
PropagationController | /api/v1/lineage/propagation | POST /run, GET /inherited/{datasetId} |
ExportController | /api/v1/lineage/export | eksport pochodzenia danych |
monitor-service
monitor-service obsługuje monitoring potoków i systemu: alertowanie, śledzenie kosztów, wskaźnik wypalania SLA/SLO, świeżość danych, monitorowanie zmian schematu, zaplanowane raporty, routing powiadomień i skrzynkę odbiorczą, wykrywanie anomalii, playbooki oraz strumienie SSE czasu rzeczywistego. Jest oznaczony @EnableScheduling i instrumentowany OpenTelemetry + Prometheus. Klasa aplikacji: MonitorServiceApplication.kt (skanuje monitor + common).
Port: SERVER_PORT (domyślnie 8080, host 8084). Flyway db/migration/monitor, historia flyway_schema_history_monitor; eksport Prometheus włączony; Jackson SNAKE_CASE.
Kluczowe klasy
- Alerty —
service/AlertService.kt,AlertDefinitionService.kt,alert/AlertDefinition.kt,alert/NotificationService.kt. - Metryki —
service/MetricsService.kt,SystemMetricsCollector.kt,KubernetesMetricsService.kt,KeycloakMetricsService.kt,ClusterMetricsHistoryService.kt,PerformanceMetricsService.kt. - Koszty —
service/CostMetricsService.kt,CostHistoryService.kt,CostAnomalyDetector.kt,PipelineCostService.kt, zconfig/GcpCostProperties.kt. - SLA / świeżość / anomalie —
service/SlaBurnRateService.kt,FreshnessTracker.kt,AnomalyBaselineService.kt. - Zmiana schematu —
service/SchemaChangeMonitor.kt. - Raporty —
service/ReportGenerationService.kt,PipelineRunLogService.kt. - Powiadomienia —
notification/NotificationRouter.kt,ChannelAdapters.kt,NotificationInboxController.kt,NotificationChannelService.kt. - Playbooki —
playbook/PlaybookEngine.kt. - Czas rzeczywisty —
realtime/SseAlertController.kt,AlertEventPublisher.kt,MetricsStreamService.kt. - Konfiguracja —
OpenTelemetryConfig.kt,MetricsExporterConfig.kt,filter/TracingFilter.kt.
Powierzchnia REST
| Kontroler | Ścieżka bazowa | Kluczowe operacje |
|---|---|---|
AlertController | /api/v1/monitor/alerts | list/get/create/delete, ack/resolve, operacje masowe, zdarzenia, historia, podsumowanie, czyszczenie |
AlertDefinitionController | /api/v1/monitor | CRUD definicji + toggle/mute, CRUD kanałów powiadomień + test |
DashboardController | /api/v1/monitor | /dashboard, /pipelines/{id}/metrics, /connectors/health, /system/health, POST /pipelines/run |
Rodzina CostController | /api/v1/monitor/costs | podsumowanie, by-pipeline, by-workspace, anomalie, trend, CRUD budżetów, stawki |
FreshnessController | /api/v1/monitor/freshness | CRUD reguł, historia, ewaluacja |
SlaBurnRateController | /api/v1/monitor/sla | potoki, CRUD definicji |
PerformanceMetricsController | /api/v1/monitor/metrics/performance | percentyle, przepustowość, porównanie, podsumowanie |
SchemaChangeController | /api/v1/monitor | schema-changes ack/resolve, schema-monitor register/scan, migawki |
AnomalyBaselineController | /api/v1/monitor/anomalies | poziomy bazowe, detect, recompute |
ReportController | /api/v1/reports | CRUD harmonogramów + toggle/send-now, podgląd, typy |
NotificationInboxController | /api/v1/notifications | list, unread-count, mark-read, GET /stream (SSE) |
PlaybookController | /api/v1/monitor/playbooks | list/create/delete, match |
SseAlertController | /api/v1/monitor/sse | /alerts, /metrics (SSE text/event-stream) |
Model danych i migracje (db/migration/monitor, V9–V25)
Encje obejmują monitor_alert_definitions/alerts/alert_events/alert_silences, data_freshness_rules/checks, metric_baselines, monitor_connector_health, monitor_pipeline_metrics/runs, pipeline_cost_records/budgets/anomalies, pipeline_run_logs, monitor_sla_burn_rate_snapshots, monitor_slo_definitions, notifications_inbox, notification_routes, playbook_routes oraz monitor_schema_change_events/snapshots. Najważniejsze migracje: V9 starsze alerty, V11 świeżość, V12 poziomy bazowe, V13 monitorowanie zmian schematu, V14 harmonogramy raportów, V15 koszty potoków, V16 wskaźnik wypalania SLA, V18 historia metryk klastra, V19 kwarantanna, V22 logi uruchomień, V24 trasy playbooków, V25 routing powiadomień/skrzynka odbiorcza. Z usługą dostarczany jest plik resources/prometheus/alert-rules.yml.
Moduły biblioteczne
connector-sdk
connector-sdk to biblioteka frameworka konektorów — bez aplikacji Spring, pakowana jako JAR konsumowany przez metadata-service i pipeline-engine. Definiuje kontrakt konektora i dostarcza 21 implementacji konektorów. Nie ma application.yml, jedynie rejestrację ServiceLoader w META-INF.
Kluczowe klasy
ConnectorSDK.kt— rdzeniowy interfejsAutoCloseable:discoverSchema(config),extractData(config): DataStream,loadData(...),validateConnection/testConnection,capabilities(), plusconnectorType/connectorName/displayName.registry/ConnectorRegistry.kt— singletonowy rejestr:register,getConnector,listConnectors;discoverConnectors()poprzezServiceLoaderJavy. Sparowany zConnectorAutoConfiguration.kt.- Klasy bazowe —
base/JDBCConnectorBase.kt,base/FileConnectorBase.kt,base/NativeConnectorBase.kt. - Model —
ConnectorConfig,ConnectorCapability(enum: SCHEMA_DISCOVERY, READ, WRITE, COLUMN_PROJECTION, PREDICATE_PUSHDOWN, SCHEMA_EVOLUTION, STREAMING, PARTITIONING, COMPRESSION, CLOUD_STORAGE, BATCH_WRITE),Schema,Table,Column,DataStream,DataBatch,DataRecord,EventStream,StorageType,ConnectorException.
Implementacje konektorów (impl/, 21)
| Kategoria | Konektory |
|---|---|
| Bazy danych (JDBC) | postgresql, mysql, mssql, oracle, sybase, teradata, hana (SAP HANA) |
| Hurtownie w chmurze | bigquery, databricks, snowflake |
| Pliki | csv, json, xml, excel, parquet (z FileFormatDetector, StorageAdapter) |
| Strumieniowanie / komunikaty | kafka (offset manager, schema registry, SerDe), pubsub, eventhubs |
| SaaS / ERP | salesforce, servicenow, sap (SapErpConnector, SapODataParser), rest |
| CDC | DebeziumCdcConnector, DebeziumCdcManager, MysqlBinlogCDC, MssqlChangeDataCapture, MongodbChangeStream |
| CDR telekomunikacyjny | CdrConnector, Asn1Decoder, CdrBinaryDecoder, CdrParquetConverter, z formatami roamingowymi Tap3Template, RapTemplate, NrtrdeTemplate |
| NoSQL | mongodb (konektor, agregacja, change stream, mapper schematu) |
pushdown-sql
pushdown-sql to usługa transpilacji dialektów SQL — tłumaczy ANSI/kanoniczny SQL na dialekty docelowych baz danych, aby logika zapytań mogła być zepchnięta (pushdown) do źródłowych lub docelowych baz danych. PushdownSqlApplication.kt to aplikacja Spring Boot, która wyłącza autokonfigurację DataSource, JPA i Security — jest to czysta usługa obliczeniowa.
Port: 8083. Brak bazy danych, brak bezpieczeństwa, brak JPA. Actuator udostępnia health/info/metrics. Brak kontrolera REST — udostępniana jako wstrzykiwalna usługa/biblioteka dla innych modułów.
Kluczowe klasy
transpiler/SqlTranspiler.kt—@Service;transpile(sql, targetDialect): TranspileResultparsuje SQL za pomocą JSqlParser (CCJSqlParserUtil), waliduje składnię, deleguje do dialektu i zbiera ostrzeżenia (nieobsługiwany join LATERAL, typ ARRAY, nadmiernie długie identyfikatory). DostarczatranslateFunction,translateDataType,getAvailableDialects.transpiler/dialect/Dialect.kt— interfejs dialektu:name,transpile, cytowanie identyfikatorów, flagi możliwości (supportsWindowFunctions,supportsCTE,supportsLateralJoin,supportsArrayType),maxIdentifierLength,translateFunction,translateDataType,translateLimitOffset/paginationSql,dualTable,castExpression.- Implementacje dialektów (8) —
PostgresDialect,MssqlDialect,OracleDialect,SnowflakeDialect,DatabricksDialect,HanaDialect,SybaseDialect,TeradataDialect. - Wizytatory —
visitor/SqlAstVisitor.kt,FunctionMapper.kt,DataTypeMapper.kt. TranspileResult— flaga sukcesu, oryginalny/transpilowany SQL, dialekt docelowy, ostrzeżenia, błąd.
Uwagi przekrojowe
- Propagacja tożsamości — brama waliduje JWT Keycloak raz; usługi Spring MVC w dół stosu ufają nagłówkom
X-User-*wstrzykniętym przez bramę (odczytywanym doSecurityContextHolderzcommon). Każda usługa niezależnie konfiguruje także serwer zasobów OAuth2 jako głęboką obronę. - Strategia bazy danych — metadata-service jest właścicielem schematu katalogu/nadzoru; lineage-service współdzieli tę bazę danych z wyłączonym Flyway; pipeline-engine i monitor-service utrzymują własne zestawy migracji zarządzanych przez Flyway z per-usługowymi tabelami historii.
- Czas rzeczywisty — SSE dla alertów/metryk monitora i skrzynki powiadomień; WebSocket do strumieniowania statusu uruchomień potoków.
- Specyfika telekomunikacyjna — konektor CDR z szablonami rekordów roamingowych ASN.1 / TAP3 / RAP / NRTRDE, kontroler metryk CDR w metadata-service, moduły RODO i zgodności z polskim prawem oraz harmonogramy domyślnie ustawione na
Europe/Warsaw.