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żkaOperation IDPrzeznaczenieUwierzytelnianie
GET/api/v1/pipelineslistPipelinesLista pipeline'ówJWT
POST/api/v1/pipelinescreatePipelineUtworzenie pipeline'u (zaczyna w DRAFT)JWT
GET/api/v1/pipelines/{id}getPipelinePobranie pipeline'u wraz z YAMLJWT
PUT/api/v1/pipelines/{id}updatePipelineAktualizacja pipeline'u (automatycznie inkrementuje wersję)JWT
DELETE/api/v1/pipelines/{id}deletePipelineUsunięcie pipeline'u + wszystkich wersji/uruchomieńJWT

Parametry

Punkt końcowyParametrWTypUwagi
listPipelinesworkspaceIdqueryuuidFiltrowanie po workspace
listPipelinessearchquerystringFiltr nazwy z tekstem dowolnym
getPipeline / updatePipeline / deletePipelineidpathuuidIdentyfikator 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żkaOperation IDPrzeznaczenieUwierzytelnianie
POST/api/v1/pipelines/{id}/runrunPipelineKompilacja + zaplanowanie wykonania; zwraca natychmiastJWT
GET/api/v1/runs/{runId}getRunStatusPobranie statusu uruchomienia + statusów per-zadanieJWT
POST/api/v1/runs/{runId}/cancelcancelRunAnulowanie trwającego wykonaniaJWT

Parametry i ciała

Punkt końcowyParametrWTypUwagi
runPipelineidpathuuidIdentyfikator pipeline'u
runPipelinebodyRunRequestOpcjonalne yamlDefinition, triggeredBy
getRunStatus / cancelRunrunIdpathuuidIdentyfikator 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żkaOperation IDPrzeznaczenieUwierzytelnianie
GET/api/v1/scheduler/scheduleslistSchedulesLista wszystkich harmonogramówJWT
POST/api/v1/scheduler/schedulescreateScheduleUtworzenie harmonogramu cronJWT
GET/api/v1/scheduler/schedules/{id}getSchedulePobranie harmonogramuJWT
PUT/api/v1/scheduler/schedules/{id}updateScheduleAktualizacja harmonogramuJWT
DELETE/api/v1/scheduler/schedules/{id}deleteScheduleUsunięcie harmonogramuJWT
POST/api/v1/scheduler/schedules/{id}/pausepauseScheduleWstrzymanie harmonogramuJWT
POST/api/v1/scheduler/schedules/{id}/resumeresumeScheduleWznowienie wstrzymanego harmonogramuJWT
GET/api/v1/scheduler/next-runsgetNextRunsPodglą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żkaOperation IDPrzeznaczenieUwierzytelnianie
POST/api/v1/scheduler/backfillcreateBackfillUtworzenie zadania backfill dla zakresu datJWT
GET/api/v1/scheduler/backfill/{id}getBackfillProgressPobranie postępu backfilluJWT
POST/api/v1/scheduler/backfill/{id}/pausepauseBackfillWstrzymanie backfilluJWT
POST/api/v1/scheduler/backfill/{id}/resumeresumeBackfillWznowienie backfilluJWT
POST/api/v1/scheduler/backfill/{id}/cancelcancelBackfillAnulowanie backfilluJWT

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żkaOperation IDPrzeznaczenieUwierzytelnianie
GET/api/v1/scheduler/dependencies/{pipelineId}getDependencyGraphPobranie grafu zależności w górę/w dółJWT
POST/api/v1/scheduler/dependenciesaddDependencyUtworzenie zależności między pipeline'amiJWT
DELETE/api/v1/scheduler/dependencies/{id}removeDependencyUsunięcie zależnościJWT

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żkaOperation IDPrzeznaczenieUwierzytelnianie
POST/api/v1/orchestrator/triggertriggerPipelineZewnętrzny wyzwalacz webhookHMAC
POST/api/v1/orchestrator/callbackreceiveCallbackOdebranie callbacku ukończenia/niepowodzenia
GET/api/v1/orchestrator/airflow/dagslistSyncedDagsLista DAG-ów wygenerowanych z pipeline'ówJWT
POST/api/v1/orchestrator/airflow/sync/{pipelineId}syncPipelineAsAirflowDagWygenerowanie + synchronizacja DAG-a Airflow z YAMLJWT
POST/api/v1/orchestrator/airflow/trigger/{dagId}triggerAirflowDagWyzwolenie uruchomienia DAG-a AirflowJWT
GET/api/v1/orchestrator/airflow/status/{dagId}/{dagRunId}getAirflowDagRunStatusPobranie statusu uruchomienia DAG-a AirflowJWT
GET/api/v1/orchestrator/webhookslistWebhooksLista zarejestrowanych webhookówJWT
POST/api/v1/orchestrator/webhooksregisterWebhookRejestracja punktu końcowego webhookJWT

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żkaOperation IDPrzeznaczenieUwierzytelnianie
GET/api/v1/git/reposlistReposLista workspace'ów z repozytoriami GitJWT
POST/api/v1/git/reposconnectRepoSklonowanie zdalnego / inicjalizacja lokalnego repozytoriumJWT
GET/api/v1/git/history/{pipelineId}getHistoryPobranie historii wersji pipeline'uJWT
GET/api/v1/git/diff/{pipelineId}diffVersionsPorównanie dwóch wersji pipeline'uJWT
POST/api/v1/git/rollback/{pipelineId}/{version}rollbackWycofanie pipeline'u do wersjiJWT
GET/api/v1/git/brancheslistBranchesLista gałęziJWT
POST/api/v1/git/branchescreateBranchUtworzenie gałęziJWT
POST/api/v1/git/mergemergeBranchScalenie gałęziJWT
POST/api/v1/git/webhooksreceiveGitWebhookOdebranie zdarzeń push i PR/MR z GitHub/GitLabHMAC

Parametry

Punkt końcowyParametrWTypUwagi
getHistorypipelineIdpathuuid
getHistoryworkspaceSlugquerystringDomyślnie default
getHistorymaxVersionsqueryintDomyślnie 50
diffVersionsfrom, toquerystringWymagane
diffVersionsworkspaceSlugquerystring
rollbackpipelineId, versionpath
rollbackworkspaceSlugquerystring
listBranchesworkspaceSlugquerystring
listBranchesincludeRemotequerybool
receiveGitWebhookX-GitHub-Event, X-Hub-Signature-256, X-Gitlab-Event, X-Gitlab-TokenheaderstringHMAC 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"
}

Punkty końcowe Flink zarządzają zadaniami strumieniowymi oraz klastrami sesyjnymi na Kubernetes.

MetodaŚcieżkaOperation IDPrzeznaczenieUwierzytelnianie
POST/api/v1/flink/jobssubmitFlinkJobZłożenie zadania Flink do KubernetesJWT
GET/api/v1/flink/jobslistFlinkJobsLista aktywnych (nieterminalnych) zadańJWT
GET/api/v1/flink/jobs/{jobId}getFlinkJobStatusOdpytanie Flink REST API o status zadaniaJWT
DELETE/api/v1/flink/jobs/{jobId}removeFlinkJobUsunięcie zakończonego zadania ze śledzeniaJWT
POST/api/v1/flink/jobs/{jobId}/cancelcancelFlinkJobAnulowanie działającego zadania (likwiduje klaster aplikacji)JWT
POST/api/v1/flink/jobs/{jobId}/savepointtriggerFlinkSavepointWyzwolenie savepointu (przechowywanego w GCS)JWT
GET/api/v1/flink/jobs/{jobId}/metricsgetFlinkJobMetricsMetryki przepustowości/checkpointów/backpressureJWT
GET/api/v1/flink/jobs/{jobId}/logsgetFlinkJobLogsPobranie wierszy logu TaskManagerJWT

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"
}
MetodaŚcieżkaOperation IDPrzeznaczenieUwierzytelnianie
GET/api/v1/flink/clusterslistFlinkClustersLista klastrów sesyjnychJWT
POST/api/v1/flink/clusterscreateFlinkClusterUtworzenie klastra sesyjnego na K8sJWT
GET/api/v1/flink/clusters/{name}getFlinkClusterPobranie szczegółów klastraJWT
DELETE/api/v1/flink/clusters/{name}deleteFlinkClusterUsunięcie klastra, zwolnienie zasobów K8sJWT
GET/api/v1/flink/clusters/{name}/healthgetFlinkClusterHealthKondycja klastra, gotowość TaskManager/slotówJWT
POST/api/v1/flink/clusters/{name}/scalescaleFlinkClusterZmiana liczby TaskManagerJWT

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.

Poprzednia
Przegląd API