Dokumentacja API
Dokumentacja API: Pipeline Engine
Pipeline Engine (openapi-pipeline-engine.yaml, port 8082) jest właścicielem definicji pipeline'ów, wykonywania, planowania, orkiestracji, wersjonowania opartego na Git oraz zadań strumieniowych Flink. Wszystkie punkty końcowe są serwowane pod /api/v1 i wymagają JWT bearerAuth — punkty końcowe webhooków uwierzytelniają się zamiast tego nagłówkami sygnatury HMAC.
Pipeline'y
Definicje pipeline'ów są zarządzane przez metadata-service za prefiksem bramy /api/v1/pipelines/**. Pipeline niesie definicję YAML oraz automatycznie inkrementowaną wersję. Nowe pipeline'y zaczynają w stanie DRAFT.
| Metoda | Ścieżka | Operation ID | Przeznaczenie | Uwierzytelnianie |
|---|---|---|---|---|
| GET | /api/v1/pipelines | listPipelines | Lista pipeline'ów | JWT |
| POST | /api/v1/pipelines | createPipeline | Utworzenie pipeline'u (zaczyna w DRAFT) | JWT |
| GET | /api/v1/pipelines/{id} | getPipeline | Pobranie pipeline'u wraz z YAML | JWT |
| PUT | /api/v1/pipelines/{id} | updatePipeline | Aktualizacja pipeline'u (automatycznie inkrementuje wersję) | JWT |
| DELETE | /api/v1/pipelines/{id} | deletePipeline | Usunięcie pipeline'u + wszystkich wersji/uruchomień | JWT |
Parametry
| Punkt końcowy | Parametr | W | Typ | Uwagi |
|---|---|---|---|---|
listPipelines | workspaceId | query | uuid | Filtrowanie po workspace |
listPipelines | search | query | string | Filtr nazwy z tekstem dowolnym |
getPipeline / updatePipeline / deletePipeline | id | path | uuid | Identyfikator pipeline'u |
Model Pipeline
Kluczowe pola Pipeline: id, workspaceId, name, description, yamlDefinition, version (int ≥ 1), status, scheduleCron, tags[], sourceConnectionId, targetConnectionId, maxRetries (domyślnie 3), retryDelaySeconds (domyślnie 60), timeoutSeconds (domyślnie 3600), notifications, isActive, createdAt, updatedAt, createdBy.
Enum status: DRAFT | SCHEDULED | RUNNING | SUCCESS | FAILED | WARNING | CANCELLED.
PipelineCreate wymaga name oraz yamlDefinition.
// POST /api/v1/pipelines → 201 Created
// Request: PipelineCreate
{
"name": "billing-daily-load",
"description": "Daily Teradata → Snowflake billing extract",
"yamlDefinition": "version: 1\nsource:\n connection: teradata-prod\ntarget:\n connection: snowflake-dwh\n",
"scheduleCron": "0 2 * * *",
"tags": ["billing", "daily"],
"sourceConnectionId": "a1b2c3d4-0000-0000-0000-000000000001",
"targetConnectionId": "a1b2c3d4-0000-0000-0000-000000000002"
}
// Response: Pipeline
{
"id": "f47ac10b-58cc-4372-a567-0e02b2c3d479",
"workspaceId": "9c1f0e2a-1111-2222-3333-444455556666",
"name": "billing-daily-load",
"description": "Daily Teradata → Snowflake billing extract",
"yamlDefinition": "version: 1\nsource:\n connection: teradata-prod\n...",
"version": 1,
"status": "DRAFT",
"scheduleCron": "0 2 * * *",
"tags": ["billing", "daily"],
"maxRetries": 3,
"retryDelaySeconds": 60,
"timeoutSeconds": 3600,
"isActive": true,
"createdAt": "2025-12-15T14:30:00Z",
"updatedAt": "2025-12-15T14:30:00Z",
"createdBy": "u-1001"
}
Uwaga
updatePipeline nie modyfikuje istniejącej wersji — automatycznie inkrementuje version i zapisuje nową rewizję. Pełna historia rewizji jest dostępna przez punkty końcowe Git poniżej. deletePipeline kaskadowo usuwa każdą wersję i każde uruchomienie.
Wykonywanie i uruchomienia
Uruchomienie to skompilowane, zaplanowane wykonanie pipeline'u. runPipeline zwraca natychmiast z 202 Accepted — odpytuj getRunStatus, aby śledzić postęp.
| Metoda | Ścieżka | Operation ID | Przeznaczenie | Uwierzytelnianie |
|---|---|---|---|---|
| POST | /api/v1/pipelines/{id}/run | runPipeline | Kompilacja + zaplanowanie wykonania; zwraca natychmiast | JWT |
| GET | /api/v1/runs/{runId} | getRunStatus | Pobranie statusu uruchomienia + statusów per-zadanie | JWT |
| POST | /api/v1/runs/{runId}/cancel | cancelRun | Anulowanie trwającego wykonania | JWT |
Parametry i ciała
| Punkt końcowy | Parametr | W | Typ | Uwagi |
|---|---|---|---|---|
runPipeline | id | path | uuid | Identyfikator pipeline'u |
runPipeline | body | — | RunRequest | Opcjonalne yamlDefinition, triggeredBy |
getRunStatus / cancelRun | runId | path | uuid | Identyfikator uruchomienia |
RunResponse: runId, pipelineId, status (RUNNING | PENDING), taskCount, engine (SPARK | FLINK | JDBC). RunStatusResponse.overallStatus: PENDING | RUNNING | COMPLETED | FAILED | CANCELLED; taskStatuses oraz errors to mapy kluczowane identyfikatorem zadania.
// POST /api/v1/pipelines/{id}/run → 202 Accepted
// Request: RunRequest (all fields optional)
{
"triggeredBy": "u-1001"
}
// Response: RunResponse
{
"runId": "3fa85f64-5717-4562-b3fc-2c963f66afa6",
"pipelineId": "f47ac10b-58cc-4372-a567-0e02b2c3d479",
"status": "RUNNING",
"taskCount": 4,
"engine": "SPARK"
}
// GET /api/v1/runs/{runId} → 200 OK (RunStatusResponse)
{
"runId": "3fa85f64-5717-4562-b3fc-2c963f66afa6",
"pipelineId": "f47ac10b-58cc-4372-a567-0e02b2c3d479",
"overallStatus": "RUNNING",
"taskStatuses": {
"extract": "COMPLETED",
"transform": "RUNNING",
"load": "PENDING"
},
"errors": {}
}
// POST /api/v1/runs/{runId}/cancel → 200 OK
{
"runId": "3fa85f64-5717-4562-b3fc-2c963f66afa6",
"status": "CANCELLED",
"message": "Run cancelled by user request"
}
Uwaga
Tabela routingu bramy odwołuje się do strumienia WebSocket pod /api/v1/runs/*/stream dla aktualizacji uruchomień na żywo, ale nie jest on zdefiniowany jako formalna operacja na ścieżce w specyfikacji. Jako udokumentowany mechanizm stosuj odpytywanie getRunStatus.
Scheduler
Scheduler zarządza harmonogramami cron, zadaniami backfill oraz grafami zależności pipeline'ów.
Harmonogramy
| Metoda | Ścieżka | Operation ID | Przeznaczenie | Uwierzytelnianie |
|---|---|---|---|---|
| GET | /api/v1/scheduler/schedules | listSchedules | Lista wszystkich harmonogramów | JWT |
| POST | /api/v1/scheduler/schedules | createSchedule | Utworzenie harmonogramu cron | JWT |
| GET | /api/v1/scheduler/schedules/{id} | getSchedule | Pobranie harmonogramu | JWT |
| PUT | /api/v1/scheduler/schedules/{id} | updateSchedule | Aktualizacja harmonogramu | JWT |
| DELETE | /api/v1/scheduler/schedules/{id} | deleteSchedule | Usunięcie harmonogramu | JWT |
| POST | /api/v1/scheduler/schedules/{id}/pause | pauseSchedule | Wstrzymanie harmonogramu | JWT |
| POST | /api/v1/scheduler/schedules/{id}/resume | resumeSchedule | Wznowienie wstrzymanego harmonogramu | JWT |
| GET | /api/v1/scheduler/next-runs | getNextRuns | Podgląd nadchodzących czasów uruchomień | JWT |
Pola CreateScheduleRequest: cronExpression (5-polowe), timezone, businessDaysOnly, skipHolidays, maintenanceWindows[], overlapStrategy (SKIP | QUEUE | CANCEL_RUNNING), catchupStrategy (SKIP_TO_CURRENT | RUN_ALL | RUN_LATEST_N). getNextRuns akceptuje limit (domyślnie 20).
// POST /api/v1/scheduler/schedules → 201 Created
// Request: CreateScheduleRequest
{
"pipelineId": "f47ac10b-58cc-4372-a567-0e02b2c3d479",
"cronExpression": "0 2 * * *",
"timezone": "Europe/Warsaw",
"businessDaysOnly": false,
"skipHolidays": true,
"overlapStrategy": "SKIP",
"catchupStrategy": "SKIP_TO_CURRENT"
}
// Response: ScheduleResponse
{
"id": "7b2e1c44-aaaa-bbbb-cccc-ddddeeeeffff",
"pipelineId": "f47ac10b-58cc-4372-a567-0e02b2c3d479",
"cronExpression": "0 2 * * *",
"timezone": "Europe/Warsaw",
"status": "ACTIVE",
"nextRunAt": "2025-12-16T01:00:00Z"
}
Backfille
| Metoda | Ścieżka | Operation ID | Przeznaczenie | Uwierzytelnianie |
|---|---|---|---|---|
| POST | /api/v1/scheduler/backfill | createBackfill | Utworzenie zadania backfill dla zakresu dat | JWT |
| GET | /api/v1/scheduler/backfill/{id} | getBackfillProgress | Pobranie postępu backfillu | JWT |
| POST | /api/v1/scheduler/backfill/{id}/pause | pauseBackfill | Wstrzymanie backfillu | JWT |
| POST | /api/v1/scheduler/backfill/{id}/resume | resumeBackfill | Wznowienie backfillu | JWT |
| POST | /api/v1/scheduler/backfill/{id}/cancel | cancelBackfill | Anulowanie backfillu | JWT |
Pola CreateBackfillRequest: granularity (HOURLY | DAILY | WEEKLY | MONTHLY), concurrency, priority (LOW | MEDIUM | HIGH), dryRun, parameterOverrides.
// POST /api/v1/scheduler/backfill → 201 Created
// Request: CreateBackfillRequest
{
"pipelineId": "f47ac10b-58cc-4372-a567-0e02b2c3d479",
"from": "2025-11-01",
"to": "2025-11-30",
"granularity": "DAILY",
"concurrency": 4,
"priority": "MEDIUM",
"dryRun": false
}
// Response: BackfillResponse
{
"id": "c9d8e7f6-1234-5678-9abc-def012345678",
"pipelineId": "f47ac10b-58cc-4372-a567-0e02b2c3d479",
"status": "RUNNING",
"totalIntervals": 30,
"completedIntervals": 11,
"failedIntervals": 0
}
Zależności
| Metoda | Ścieżka | Operation ID | Przeznaczenie | Uwierzytelnianie |
|---|---|---|---|---|
| GET | /api/v1/scheduler/dependencies/{pipelineId} | getDependencyGraph | Pobranie grafu zależności w górę/w dół | JWT |
| POST | /api/v1/scheduler/dependencies | addDependency | Utworzenie zależności między pipeline'ami | JWT |
| DELETE | /api/v1/scheduler/dependencies/{id} | removeDependency | Usunięcie zależności | JWT |
AddDependencyRequest.type: STRONG | WEAK. Dodanie międzyworkspace'owej zależności zwraca 403.
// POST /api/v1/scheduler/dependencies → 201 Created
// Request: AddDependencyRequest
{
"upstreamPipelineId": "f47ac10b-58cc-4372-a567-0e02b2c3d479",
"downstreamPipelineId": "a1b2c3d4-9999-8888-7777-666655554444",
"type": "STRONG"
}
Orkiestrator
Orkiestrator integruje zewnętrzne schedulery (Airflow, AutomateNow) poprzez webhooki, callbacki oraz synchronizację DAG-ów. Webhooki wyzwalające uwierzytelniają się nagłówkami HMAC, a nie JWT.
| Metoda | Ścieżka | Operation ID | Przeznaczenie | Uwierzytelnianie |
|---|---|---|---|---|
| POST | /api/v1/orchestrator/trigger | triggerPipeline | Zewnętrzny wyzwalacz webhook | HMAC |
| POST | /api/v1/orchestrator/callback | receiveCallback | Odebranie callbacku ukończenia/niepowodzenia | — |
| GET | /api/v1/orchestrator/airflow/dags | listSyncedDags | Lista DAG-ów wygenerowanych z pipeline'ów | JWT |
| POST | /api/v1/orchestrator/airflow/sync/{pipelineId} | syncPipelineAsAirflowDag | Wygenerowanie + synchronizacja DAG-a Airflow z YAML | JWT |
| POST | /api/v1/orchestrator/airflow/trigger/{dagId} | triggerAirflowDag | Wyzwolenie uruchomienia DAG-a Airflow | JWT |
| GET | /api/v1/orchestrator/airflow/status/{dagId}/{dagRunId} | getAirflowDagRunStatus | Pobranie statusu uruchomienia DAG-a Airflow | JWT |
| GET | /api/v1/orchestrator/webhooks | listWebhooks | Lista zarejestrowanych webhooków | JWT |
| POST | /api/v1/orchestrator/webhooks | registerWebhook | Rejestracja punktu końcowego webhook | JWT |
triggerPipeline wymaga nagłówków X-Webhook-Signature oraz X-Webhook-Id, a także ciała TriggerRequest wymagającego external_job_id. WebhookRegistrationRequest wymaga name, source, hmac_secret, callback_url. triggerPipeline oraz triggerAirflowDag zwracają 202 Accepted.
// POST /api/v1/orchestrator/trigger → 202 Accepted
// Headers: X-Webhook-Signature: <hmac>, X-Webhook-Id: <uuid>
// Request: TriggerRequest
{
"external_job_id": "airflow-billing-2025-12-15",
"pipelineId": "f47ac10b-58cc-4372-a567-0e02b2c3d479",
"parameters": { "loadDate": "2025-12-15" }
}
// Response: TriggerResponse
{
"runId": "3fa85f64-5717-4562-b3fc-2c963f66afa6",
"external_job_id": "airflow-billing-2025-12-15",
"status": "ACCEPTED"
}
// POST /api/v1/orchestrator/airflow/trigger/{dagId} → 202 Accepted
{
"dagId": "dataflow_billing_daily_load",
"dagRunId": "manual__2025-12-15T14:30:00+00:00",
"state": "queued"
}
Git
Definicje pipeline'ów są wersjonowane w Git. Te punkty końcowe zarządzają repozytoriami, historią, różnicami, wycofywaniem, gałęziami oraz scalaniami. Webhook Git uwierzytelnia się nagłówkami HMAC dostawcy.
| Metoda | Ścieżka | Operation ID | Przeznaczenie | Uwierzytelnianie |
|---|---|---|---|---|
| GET | /api/v1/git/repos | listRepos | Lista workspace'ów z repozytoriami Git | JWT |
| POST | /api/v1/git/repos | connectRepo | Sklonowanie zdalnego / inicjalizacja lokalnego repozytorium | JWT |
| GET | /api/v1/git/history/{pipelineId} | getHistory | Pobranie historii wersji pipeline'u | JWT |
| GET | /api/v1/git/diff/{pipelineId} | diffVersions | Porównanie dwóch wersji pipeline'u | JWT |
| POST | /api/v1/git/rollback/{pipelineId}/{version} | rollback | Wycofanie pipeline'u do wersji | JWT |
| GET | /api/v1/git/branches | listBranches | Lista gałęzi | JWT |
| POST | /api/v1/git/branches | createBranch | Utworzenie gałęzi | JWT |
| POST | /api/v1/git/merge | mergeBranch | Scalenie gałęzi | JWT |
| POST | /api/v1/git/webhooks | receiveGitWebhook | Odebranie zdarzeń push i PR/MR z GitHub/GitLab | HMAC |
Parametry
| Punkt końcowy | Parametr | W | Typ | Uwagi |
|---|---|---|---|---|
getHistory | pipelineId | path | uuid | — |
getHistory | workspaceSlug | query | string | Domyślnie default |
getHistory | maxVersions | query | int | Domyślnie 50 |
diffVersions | from, to | query | string | Wymagane |
diffVersions | workspaceSlug | query | string | — |
rollback | pipelineId, version | path | — | — |
rollback | workspaceSlug | query | string | — |
listBranches | workspaceSlug | query | string | — |
listBranches | includeRemote | query | bool | — |
receiveGitWebhook | X-GitHub-Event, X-Hub-Signature-256, X-Gitlab-Event, X-Gitlab-Token | header | string | HMAC dostawcy |
ConnectRepoRequest.credentialType: NONE | TOKEN | SSH_KEY. WebhookResult.action: SYNC | DEPLOY | VALIDATE | IGNORE | ERROR. mergeBranch zwraca 409 w przypadku konfliktu scalania.
// POST /api/v1/git/repos → 201 Created
// Request: ConnectRepoRequest
{
"workspaceSlug": "billing-team",
"remoteUrl": "https://gitlab.polkomtel.pl/dataflow/billing.git",
"credentialType": "TOKEN",
"credentialRef": "vault://git/billing-token"
}
// GET /api/v1/git/diff/{pipelineId}?from=3&to=5 → 200 OK (DiffResponse)
{
"pipelineId": "f47ac10b-58cc-4372-a567-0e02b2c3d479",
"fromVersion": 3,
"toVersion": 5,
"diff": "@@ -2,3 +2,3 @@\n- connection: teradata-stage\n+ connection: teradata-prod"
}
// POST /api/v1/git/merge → 200 OK (VersionOperationResult)
{
"success": true,
"mergedVersion": 6,
"message": "Branch feature/new-mapping merged into main"
}
Flink
Punkty końcowe Flink zarządzają zadaniami strumieniowymi oraz klastrami sesyjnymi na Kubernetes.
Zadania Flink
| Metoda | Ścieżka | Operation ID | Przeznaczenie | Uwierzytelnianie |
|---|---|---|---|---|
| POST | /api/v1/flink/jobs | submitFlinkJob | Złożenie zadania Flink do Kubernetes | JWT |
| GET | /api/v1/flink/jobs | listFlinkJobs | Lista aktywnych (nieterminalnych) zadań | JWT |
| GET | /api/v1/flink/jobs/{jobId} | getFlinkJobStatus | Odpytanie Flink REST API o status zadania | JWT |
| DELETE | /api/v1/flink/jobs/{jobId} | removeFlinkJob | Usunięcie zakończonego zadania ze śledzenia | JWT |
| POST | /api/v1/flink/jobs/{jobId}/cancel | cancelFlinkJob | Anulowanie działającego zadania (likwiduje klaster aplikacji) | JWT |
| POST | /api/v1/flink/jobs/{jobId}/savepoint | triggerFlinkSavepoint | Wyzwolenie savepointu (przechowywanego w GCS) | JWT |
| GET | /api/v1/flink/jobs/{jobId}/metrics | getFlinkJobMetrics | Metryki przepustowości/checkpointów/backpressure | JWT |
| GET | /api/v1/flink/jobs/{jobId}/logs | getFlinkJobLogs | Pobranie wierszy logu TaskManager | JWT |
FlinkJobConfig.deploymentMode: APPLICATION | SESSION. FlinkJobRecord.status: CREATED | RUNNING | FAILING | FAILED | CANCELLING | CANCELED | FINISHED | RESTARTING | SUSPENDED. triggerFlinkSavepoint wymaga path w ciele. getFlinkJobLogs akceptuje lines (domyślnie 500).
// POST /api/v1/flink/jobs → 201 Created
// Request: FlinkJobConfig
{
"jobName": "cdr-stream-enrich",
"deploymentMode": "APPLICATION",
"jarUri": "gs://dataflow-artifacts/cdr-enrich-1.4.0.jar",
"parallelism": 8
}
// Response: FlinkJobHandle
{
"jobId": "00000000-cdr-enrich-0001",
"status": "CREATED",
"deploymentMode": "APPLICATION"
}
// POST /api/v1/flink/jobs/{jobId}/savepoint → 200 OK
// Request: { "path": "gs://dataflow-savepoints/cdr-enrich/" }
{
"savepointPath": "gs://dataflow-savepoints/cdr-enrich/savepoint-abc123"
}
Klastry Flink
| Metoda | Ścieżka | Operation ID | Przeznaczenie | Uwierzytelnianie |
|---|---|---|---|---|
| GET | /api/v1/flink/clusters | listFlinkClusters | Lista klastrów sesyjnych | JWT |
| POST | /api/v1/flink/clusters | createFlinkCluster | Utworzenie klastra sesyjnego na K8s | JWT |
| GET | /api/v1/flink/clusters/{name} | getFlinkCluster | Pobranie szczegółów klastra | JWT |
| DELETE | /api/v1/flink/clusters/{name} | deleteFlinkCluster | Usunięcie klastra, zwolnienie zasobów K8s | JWT |
| GET | /api/v1/flink/clusters/{name}/health | getFlinkClusterHealth | Kondycja klastra, gotowość TaskManager/slotów | JWT |
| POST | /api/v1/flink/clusters/{name}/scale | scaleFlinkCluster | Zmiana liczby TaskManager | JWT |
FlinkClusterInfo.status: STARTING | RUNNING | SCALING | DEGRADED | STOPPING | STOPPED | ERROR. scaleFlinkCluster wymaga taskManagers (≥ 1) w ciele; createFlinkCluster zwraca 409, jeśli klaster o tej nazwie już istnieje.
// POST /api/v1/flink/clusters → 201 Created
// Request: FlinkClusterCreateRequest
{
"name": "streaming-prod",
"taskManagers": 3,
"taskSlotsPerManager": 4
}
// POST /api/v1/flink/clusters/{name}/scale → 200 OK (FlinkClusterInfo)
// Request: { "taskManagers": 6 }
{
"name": "streaming-prod",
"status": "SCALING",
"taskManagers": 6,
"taskSlotsPerManager": 4
}
Uwaga
Anulowanie zadania w trybie APPLICATION likwiduje jego dedykowany klaster aplikacji. W przypadku długotrwałych zadań strumieniowych najpierw wyzwól savepoint, aby zadanie mogło zostać wznowione ze spójnego checkpointu przechowywanego w GCS.