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
| Method | Endpoint | Description |
|---|---|---|
POST | /api/v1/jobs | Submit new job |
GET | /api/v1/jobs | List 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}/savepoints | Create savepoint |
GET | /api/v1/jobs/{job_id}/metrics | Get job metrics |
GET | /api/v1/jobs/{job_id}/logs | Get job logs |
POST | /api/v1/jobs/{job_id}/restart | Restart 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 jobsheliosdb_job_status(gauge): Job status by job_idheliosdb_job_events_processed_total(counter): Total events processedheliosdb_job_throughput(gauge): Current throughput (events/sec)heliosdb_job_latency_seconds(histogram): Processing latencyheliosdb_job_checkpoint_duration_seconds(histogram): Checkpoint durationheliosdb_job_checkpoint_failures_total(counter): Checkpoint failuresheliosdb_job_backpressure_events_total(counter): Backpressure eventsheliosdb_job_restarts_total(counter): Job restartsheliosdb_job_errors_total(counter): Job errors
Document Version: 1.0