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łPortTypDB / FlywayPrzeznaczenie
api-gateway8080Spring Cloud Gateway (reaktywny WebFlux)brak (Redis)Routing, walidacja JWT, ograniczanie szybkości, maskowanie PII
metadata-service8081Spring MVC (JPA)własna DB, Flyway metadata (V1–V54)Katalog, nadzór, RODO, jakość, metadane potoków
pipeline-engine8082Spring MVC (JPA)własna DB, Flyway engine (V1–V9)Wykonywanie potoków, harmonogramowanie, samonaprawa
pushdown-sql8083Spring Boot (bez DB/bezpieczeństwa)brakBiblioteka/usługa transpilacji dialektów SQL
lineage-service8084Spring MVC (JPA)odczytuje DB metadata, Flyway wyłączonyŚledzenie pochodzenia danych i graf
monitor-service8080 (SERVER_PORT)Spring MVC (JPA)własna DB, Flyway monitor (V9–V25)Monitoring, alerty, koszty, SLA, powiadomienia
connector-sdknd.Moduł bibliotecznybrakFramework konektorów + 21 implementacji konektorów
commonnd.Moduł bibliotecznybrakWspół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

KlasaRola
config/RouteConfig.ktDeklaratywne 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.ktKonwertuje roszczenia JWT Keycloak na uprawnienia Spring.
filter/AuthFilter.ktGlobalFilter (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.ktLimiter 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.ktMaskowanie 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. Enum DataFlowRole: ADMIN(100) > ENGINEER(75) > ANALYST(50) > STEWARD(40) > VIEWER(25). Enum Permission z requiredRole (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ę odbiorcy aud (domyślnie dataflow-api), zwraca AuthenticatedPrincipal.
  • security/KeycloakJwtConverter.kt — wyodrębnia nazwę principal + uprawnienia z JWT (wariant servletowy).
  • security/SecurityContext.kt — klasa danych SecurityContext (userId, email, role, workspaceId, groups) z isAdmin / isEngineerOrAbove / isAnalystOrAbove; ThreadLocal SecurityContextHolder wypeł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)

WersjaZawartość
V1–V5Rdzeń schematu, potoki, monitoring, pochodzenie danych, audyt
V9–V14Katalog, glosariusz, rejestr schematów, klasyfikacja, nadzór, przepływ pracy
V16–V21Kontrakty 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 / V35Katalog marketplace konektorów, widok zmaterializowany wyszukiwania globalnego
V44–V54Ramy 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

  • Wykonywanieexecution/PipelineRunner.kt orkiestruje 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 poprzez ExecutionContext. Wspierany przez TaskExecutor.kt oraz ExecutionContext.kt.
  • DAGdag/DagBuilder.kt, dag/ExecutionDAG.kt, dag/DagNode.kt, z wykrywaniem cykli i zwisających zależności.
  • DSLdsl/PipelineYamlParser.kt, dsl/PipelineYamlSchema.kt.
  • Walidacjavalidation/PipelineValidator.kt, validation/ParameterResolver.kt (zmienne wbudowane, zmienne środowiskowe, parametry uruchomieniowe).
  • Warstwa usługservice/PipelineCompiler.kt, ExecutionPlanner.kt, TaskScheduler.kt, PipelineRunLogPublisher.kt.
  • Harmonogramowaniescheduler/CronScheduler.kt, BackfillManager.kt, DependencyChain.kt, RetryPolicy.kt.
  • Orkiestratoryorchestrator/AirflowClient.kt, AirflowDagGenerator.kt, AutomateNowAdapter.kt, WebhookCallbackService.kt.
  • Silniki obliczenioweflink/ (FlinkClusterManager, FlinkJobBuilder/Submitter, FlinkCheckpointManager, FlinkSqlBridge) oraz spark/ (SparkClusterManager, SparkJobBuilder, DataprocJobSubmitter, SparkSqlBridge).
  • Samonaprawahealing/SelfHealingService.kt, FailureClassifier.kt, RecoveryStrategies.kt.
  • Gitgit/GitRepositoryManager.kt, PipelineVersionControl.kt, YamlDiffEngine.kt, GitWebhookHandler.kt.
  • Usuwanie (RODO)erasure/ErasureOrchestrator.kt, ErasureController.
  • Jakośćquality/DataQuarantineService.kt, streaming/StreamingQualityEngine.kt.
  • Czas rzeczywistyrealtime/PipelineStatusWebSocket.kt, ExecutionEventPublisher.kt.
  • Pochodzenie danychlineage/OpenLineageEmitter.kt.

Powierzchnia REST

KontrolerŚcieżka bazowaKluczowe operacje
ExecutionController/api/v1POST /pipelines/{id}/run, GET /runs/{runId}, POST /runs/{runId}/cancel, POST /runs/{runId}/status
SchedulerController/api/v1/schedulerCRUD 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/gitrepozytoria, historia, diff, rollback, gałęzie, merge, webhooki
HealingController/api/v1/pipelines/{id}/healing-history, /healing-summary, /healing-policy
QuarantineController/api/v1/quarantinelist/get, approve/reject/edit-approve, operacje masowe, reguły auto-zwolnienia
ErasureController/api/v1/dsarPOST /{dsarId}/erasure/orchestrate, status uruchomienia usuwania
StreamingQualityController/api/v1/monitor/streamingjakość per-potok, reguły, reset, SQL
PipelineReportController/api/v1/pipelines/{id}/reportraporty 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 najpierw LineageGraphBuilder, wracając do magazynu w pamięci ConcurrentHashMap plus LineagePersistence.
  • 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.ktDatasetLineageRepository, ColumnLineageRepository, ponownie wykorzystujące encje metadata-service; obsługuje findLineageAsOf (podróż w czasie) oraz filtrowanie PII.
  • model/LineageEvent.ktRunEvent w stylu OpenLineage, LineageGraph, LineageNode, LineageJobNode, SchemaField.

Powierzchnia REST

KontrolerŚcieżka bazowaKluczowe operacje
LineageController/api/v1/lineageGET /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/openlineagePOST /events, GET /events
BiLineageController/api/v1/lineage/biPOST /ingest, POST /{tool}/ingest
LineageAuthoringController/api/v1/lineage/manual-edgesGET, POST, PATCH /{id}, DELETE /{id}
PropagationController/api/v1/lineage/propagationPOST /run, GET /inherited/{datasetId}
ExportController/api/v1/lineage/exporteksport 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

  • Alertyservice/AlertService.kt, AlertDefinitionService.kt, alert/AlertDefinition.kt, alert/NotificationService.kt.
  • Metrykiservice/MetricsService.kt, SystemMetricsCollector.kt, KubernetesMetricsService.kt, KeycloakMetricsService.kt, ClusterMetricsHistoryService.kt, PerformanceMetricsService.kt.
  • Kosztyservice/CostMetricsService.kt, CostHistoryService.kt, CostAnomalyDetector.kt, PipelineCostService.kt, z config/GcpCostProperties.kt.
  • SLA / świeżość / anomalieservice/SlaBurnRateService.kt, FreshnessTracker.kt, AnomalyBaselineService.kt.
  • Zmiana schematuservice/SchemaChangeMonitor.kt.
  • Raportyservice/ReportGenerationService.kt, PipelineRunLogService.kt.
  • Powiadomienianotification/NotificationRouter.kt, ChannelAdapters.kt, NotificationInboxController.kt, NotificationChannelService.kt.
  • Playbookiplaybook/PlaybookEngine.kt.
  • Czas rzeczywistyrealtime/SseAlertController.kt, AlertEventPublisher.kt, MetricsStreamService.kt.
  • KonfiguracjaOpenTelemetryConfig.kt, MetricsExporterConfig.kt, filter/TracingFilter.kt.

Powierzchnia REST

KontrolerŚcieżka bazowaKluczowe operacje
AlertController/api/v1/monitor/alertslist/get/create/delete, ack/resolve, operacje masowe, zdarzenia, historia, podsumowanie, czyszczenie
AlertDefinitionController/api/v1/monitorCRUD 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/costspodsumowanie, by-pipeline, by-workspace, anomalie, trend, CRUD budżetów, stawki
FreshnessController/api/v1/monitor/freshnessCRUD reguł, historia, ewaluacja
SlaBurnRateController/api/v1/monitor/slapotoki, CRUD definicji
PerformanceMetricsController/api/v1/monitor/metrics/performancepercentyle, przepustowość, porównanie, podsumowanie
SchemaChangeController/api/v1/monitorschema-changes ack/resolve, schema-monitor register/scan, migawki
AnomalyBaselineController/api/v1/monitor/anomaliespoziomy bazowe, detect, recompute
ReportController/api/v1/reportsCRUD harmonogramów + toggle/send-now, podgląd, typy
NotificationInboxController/api/v1/notificationslist, unread-count, mark-read, GET /stream (SSE)
PlaybookController/api/v1/monitor/playbookslist/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 interfejs AutoCloseable: discoverSchema(config), extractData(config): DataStream, loadData(...), validateConnection/testConnection, capabilities(), plus connectorType / connectorName / displayName.
  • registry/ConnectorRegistry.kt — singletonowy rejestr: register, getConnector, listConnectors; discoverConnectors() poprzez ServiceLoader Javy. Sparowany z ConnectorAutoConfiguration.kt.
  • Klasy bazowebase/JDBCConnectorBase.kt, base/FileConnectorBase.kt, base/NativeConnectorBase.kt.
  • ModelConnectorConfig, 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)

KategoriaKonektory
Bazy danych (JDBC)postgresql, mysql, mssql, oracle, sybase, teradata, hana (SAP HANA)
Hurtownie w chmurzebigquery, databricks, snowflake
Plikicsv, json, xml, excel, parquet (z FileFormatDetector, StorageAdapter)
Strumieniowanie / komunikatykafka (offset manager, schema registry, SerDe), pubsub, eventhubs
SaaS / ERPsalesforce, servicenow, sap (SapErpConnector, SapODataParser), rest
CDCDebeziumCdcConnector, DebeziumCdcManager, MysqlBinlogCDC, MssqlChangeDataCapture, MongodbChangeStream
CDR telekomunikacyjnyCdrConnector, Asn1Decoder, CdrBinaryDecoder, CdrParquetConverter, z formatami roamingowymi Tap3Template, RapTemplate, NrtrdeTemplate
NoSQLmongodb (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): TranspileResult parsuje 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). Dostarcza translateFunction, 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.
  • Wizytatoryvisitor/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 do SecurityContextHolder z common). 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.
Poprzednia
Architektura systemu
Następna
Usługi AI