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))
}

💻 Run locally

Copy the code above and run it on your machine

© 2026 ByteLearn.dev. Free courses for developers. · Privacy