Queue Depth Exporter
A custom Prometheus exporter that implements the prometheus.Collector interface to scrape metrics from an external service API. It translates queue depth, job counts, latency, and worker status into Prometheus-compatible metrics at scrape time.
Input: An external service endpoint that exposes a /stats JSON API.
Output: Prometheus metrics on :9101/metrics (queue depth, processed jobs, average latency, and active workers).
package main
import (
"encoding/json"
"log"
"net/http"
"time"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promhttp"
)
// ServiceStats represents the response from our external service
type ServiceStats struct {
QueueDepth int `json:"queue_depth"`
ProcessedJobs int `json:"processed_jobs"`
AvgLatencyMs float64 `json:"avg_latency_ms"`
WorkersActive int `json:"workers_active"`
}
type ServiceCollector struct {
endpoint string
client *http.Client
queueDepth *prometheus.Desc
jobsTotal *prometheus.Desc
latency *prometheus.Desc
workersUp *prometheus.Desc
scrapeErrors *prometheus.Desc
}
func NewServiceCollector(endpoint string) *ServiceCollector {
return &ServiceCollector{
endpoint: endpoint,
client: &http.Client{Timeout: 5 * time.Second},
queueDepth: prometheus.NewDesc(
"service_queue_depth",
"Number of items waiting in the queue.",
nil, nil,
),
jobsTotal: prometheus.NewDesc(
"service_processed_jobs_total",
"Total number of processed jobs.",
nil, nil,
),
latency: prometheus.NewDesc(
"service_avg_latency_seconds",
"Average processing latency in seconds.",
nil, nil,
),
workersUp: prometheus.NewDesc(
"service_workers_active",
"Number of active workers.",
nil, nil,
),
scrapeErrors: prometheus.NewDesc(
"service_scrape_errors_total",
"Total scrape errors.",
nil, nil,
),
}
}
func (c *ServiceCollector) Describe(ch chan<- *prometheus.Desc) {
ch <- c.queueDepth
ch <- c.jobsTotal
ch <- c.latency
ch <- c.workersUp
ch <- c.scrapeErrors
}
func (c *ServiceCollector) Collect(ch chan<- prometheus.Metric) {
resp, err := c.client.Get(c.endpoint + "/stats")
if err != nil {
ch <- prometheus.MustNewConstMetric(c.scrapeErrors, prometheus.CounterValue, 1)
return
}
defer resp.Body.Close()
var stats ServiceStats
if err := json.NewDecoder(resp.Body).Decode(&stats); err != nil {
ch <- prometheus.MustNewConstMetric(c.scrapeErrors, prometheus.CounterValue, 1)
return
}
ch <- prometheus.MustNewConstMetric(c.queueDepth, prometheus.GaugeValue, float64(stats.QueueDepth))
ch <- prometheus.MustNewConstMetric(c.jobsTotal, prometheus.CounterValue, float64(stats.ProcessedJobs))
ch <- prometheus.MustNewConstMetric(c.latency, prometheus.GaugeValue, stats.AvgLatencyMs/1000.0)
ch <- prometheus.MustNewConstMetric(c.workersUp, prometheus.GaugeValue, float64(stats.WorkersActive))
}
func main() {
collector := NewServiceCollector("http://localhost:9000")
prometheus.MustRegister(collector)
http.Handle("/metrics", promhttp.Handler())
log.Println("Exporter listening on :9101")
log.Fatal(http.ListenAndServe(":9101", nil))
}