Skip to content

Job Management API Documentation

Job Management API Documentation

Version: 1.0


Overview

The Job Management API provides complete control over the streaming job lifecycle: submission, monitoring, cancellation, savepoint management, and recovery. The API is exposed over HTTP REST, with Prometheus/Grafana integration for monitoring.

Key Features

  • Job submission with validation
  • Job lifecycle management (submit, run, pause, cancel, restart)
  • Savepoint management (create, restore, delete)
  • Real-time monitoring (metrics, status, logs)
  • Resource management (CPU, memory, parallelism)
  • Failure recovery (automatic restart, checkpointing)

API Overview

Endpoints

MethodEndpointDescription
POST/api/v1/jobsSubmit new job
GET/api/v1/jobsList all jobs
GET/api/v1/jobs/{job_id}Get job details
DELETE/api/v1/jobs/{job_id}Cancel job
POST/api/v1/jobs/{job_id}/savepointsCreate savepoint
GET/api/v1/jobs/{job_id}/metricsGet job metrics
GET/api/v1/jobs/{job_id}/logsGet job logs
POST/api/v1/jobs/{job_id}/restartRestart job

API Reference

1. Submit Job

Submit a new streaming job.

Endpoint: POST /api/v1/jobs

Request Body:

{
"name": "fraud_detection_pipeline",
"config": {
"parallelism": 4,
"checkpoint_interval_ms": 60000,
"state_backend": "rocksdb",
"restart_strategy": "fixed_delay"
},
"sources": [
{
"type": "kafka",
"config": {
"bootstrap_servers": ["localhost:9092"],
"topic": "transactions",
"group_id": "fraud_detector"
}
}
],
"transformations": [
{
"type": "filter",
"config": {
"condition": "amount > 10000"
}
},
{
"type": "cep",
"config": {
"pattern": "A B+ C",
"window": "1h"
}
}
],
"sinks": [
{
"type": "database",
"config": {
"connection_string": "postgresql://...",
"table": "fraud_alerts"
}
}
]
}

Response (201 Created):

{
"job_id": "job-123e4567-e89b-12d3-a456-426614174000",
"status": "submitted",
"submitted_at": "2025-10-29T10:00:00Z",
"message": "Job submitted successfully"
}

2. List Jobs

Get list of all jobs.

Endpoint: GET /api/v1/jobs

Query Parameters:

  • status (optional): Filter by status (running, finished, failed, cancelled)
  • limit (optional): Max results (default: 100)
  • offset (optional): Pagination offset

Response (200 OK):

{
"jobs": [
{
"job_id": "job-123...",
"name": "fraud_detection_pipeline",
"status": "running",
"submitted_at": "2025-10-29T10:00:00Z",
"started_at": "2025-10-29T10:00:05Z",
"uptime_seconds": 3600,
"parallelism": 4
},
{
"job_id": "job-456...",
"name": "recommendation_engine",
"status": "finished",
"submitted_at": "2025-10-29T08:00:00Z",
"started_at": "2025-10-29T08:00:03Z",
"finished_at": "2025-10-29T09:30:00Z",
"parallelism": 8
}
],
"total": 2,
"limit": 100,
"offset": 0
}

3. Get Job Details

Get detailed information about a specific job.

Endpoint: GET /api/v1/jobs/{job_id}

Response (200 OK):

{
"job_id": "job-123...",
"name": "fraud_detection_pipeline",
"status": "running",
"config": {
"parallelism": 4,
"checkpoint_interval_ms": 60000,
"state_backend": "rocksdb",
"restart_strategy": "fixed_delay"
},
"submitted_at": "2025-10-29T10:00:00Z",
"started_at": "2025-10-29T10:00:05Z",
"uptime_seconds": 3600,
"metrics": {
"events_processed": 1523400,
"throughput_per_sec": 423,
"latency_p50_ms": 3.2,
"latency_p99_ms": 12.5,
"checkpoint_count": 60,
"last_checkpoint_duration_ms": 45,
"backpressure_events": 0,
"restarts": 0
},
"tasks": [
{
"task_id": "task-1",
"name": "kafka-source",
"status": "running",
"parallelism": 1
},
{
"task_id": "task-2",
"name": "filter",
"status": "running",
"parallelism": 4
},
{
"task_id": "task-3",
"name": "cep-matcher",
"status": "running",
"parallelism": 4
},
{
"task_id": "task-4",
"name": "database-sink",
"status": "running",
"parallelism": 2
}
],
"checkpoints": [
{
"checkpoint_id": "checkpoint-1",
"timestamp": "2025-10-29T10:59:00Z",
"duration_ms": 45,
"size_bytes": 1048576,
"status": "completed"
}
]
}

4. Cancel Job

Cancel a running or paused job.

Endpoint: DELETE /api/v1/jobs/{job_id}

Query Parameters:

  • savepoint (optional): Create savepoint before cancelling (true/false, default: false)

Response (200 OK):

{
"job_id": "job-123...",
"status": "cancelled",
"cancelled_at": "2025-10-29T11:00:00Z",
"savepoint_id": "savepoint-789..." // if savepoint=true
}

5. Create Savepoint

Create a savepoint for a running job (for backup or migration).

Endpoint: POST /api/v1/jobs/{job_id}/savepoints

Request Body (optional):

{
"description": "Pre-upgrade savepoint"
}

Response (201 Created):

{
"savepoint_id": "savepoint-789...",
"job_id": "job-123...",
"created_at": "2025-10-29T11:00:00Z",
"size_bytes": 10485760,
"location": "s3://checkpoints/savepoint-789...",
"description": "Pre-upgrade savepoint"
}

6. Get Job Metrics

Get real-time metrics for a job.

Endpoint: GET /api/v1/jobs/{job_id}/metrics

Query Parameters:

  • window (optional): Time window (1m, 5m, 1h, 24h, default: 5m)

Response (200 OK):

{
"job_id": "job-123...",
"timestamp": "2025-10-29T11:00:00Z",
"metrics": {
"events_processed_total": 1523400,
"throughput_per_sec": 423,
"latency": {
"p50_ms": 3.2,
"p95_ms": 8.7,
"p99_ms": 12.5,
"p999_ms": 25.3
},
"checkpoint": {
"count": 60,
"last_duration_ms": 45,
"failures": 0
},
"backpressure": {
"events_total": 0,
"current_level": 0.0
},
"resources": {
"cpu_usage_percent": 45.2,
"memory_used_mb": 512,
"memory_total_mb": 2048
},
"errors": {
"count": 0,
"rate_per_min": 0.0
}
},
"history": [
{
"timestamp": "2025-10-29T10:55:00Z",
"throughput_per_sec": 420
},
{
"timestamp": "2025-10-29T10:56:00Z",
"throughput_per_sec": 425
}
]
}

7. Get Job Logs

Get logs for a job.

Endpoint: GET /api/v1/jobs/{job_id}/logs

Query Parameters:

  • level (optional): Filter by log level (debug, info, warn, error)
  • limit (optional): Max log lines (default: 1000)
  • tail (optional): Return last N lines (default: false)

Response (200 OK):

{
"job_id": "job-123...",
"logs": [
{
"timestamp": "2025-10-29T10:00:05.123Z",
"level": "info",
"message": "Job started successfully",
"task_id": null
},
{
"timestamp": "2025-10-29T10:00:05.456Z",
"level": "info",
"message": "Kafka source connected to localhost:9092",
"task_id": "task-1"
},
{
"timestamp": "2025-10-29T10:01:00.789Z",
"level": "info",
"message": "Checkpoint 1 completed (45ms)",
"task_id": null
},
{
"timestamp": "2025-10-29T10:05:12.345Z",
"level": "warn",
"message": "Temporary backpressure detected",
"task_id": "task-2"
}
],
"total": 4,
"limit": 1000
}

8. Restart Job

Restart a failed or cancelled job (optionally from savepoint).

Endpoint: POST /api/v1/jobs/{job_id}/restart

Request Body (optional):

{
"savepoint_id": "savepoint-789...", // Optional
"config_overrides": { // Optional
"parallelism": 8
}
}

Response (200 OK):

{
"job_id": "job-123...",
"status": "running",
"restarted_at": "2025-10-29T11:05:00Z",
"restored_from_savepoint": "savepoint-789...",
"message": "Job restarted successfully"
}

Prometheus Metrics

Job Metrics

  • heliosdb_job_count (gauge): Number of active jobs
  • heliosdb_job_status (gauge): Job status by job_id
  • heliosdb_job_events_processed_total (counter): Total events processed
  • heliosdb_job_throughput (gauge): Current throughput (events/sec)
  • heliosdb_job_latency_seconds (histogram): Processing latency
  • heliosdb_job_checkpoint_duration_seconds (histogram): Checkpoint duration
  • heliosdb_job_checkpoint_failures_total (counter): Checkpoint failures
  • heliosdb_job_backpressure_events_total (counter): Backpressure events
  • heliosdb_job_restarts_total (counter): Job restarts
  • heliosdb_job_errors_total (counter): Job errors

Document Version: 1.0