Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions cmd/blitz/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@ import (
"github.com/observiq/blitz/output/nop"
otlpgrpc "github.com/observiq/blitz/output/otlp_grpc"
promrw "github.com/observiq/blitz/output/promrw"
"github.com/observiq/blitz/output/promscrape"
stdoutout "github.com/observiq/blitz/output/stdout"
syslogout "github.com/observiq/blitz/output/syslog"
"github.com/observiq/blitz/output/tcp"
Expand Down Expand Up @@ -376,6 +377,13 @@ func run(cmd *cobra.Command, args []string) error {
logger.Error("Failed to create prometheus-remote-write output", zap.Error(err))
return err
}
case config.OutputTypePrometheusScrape:
ps := cfg.Output.PrometheusScrape
outputInstance, err = promscrape.New(ps.ListenAddress, ps.MetricsPath, ps.EmitTimestamps, ps.MetricExpiration, tel, logger)
if err != nil {
logger.Error("Failed to create prometheus-scrape output", zap.Error(err))
return err
}
default:
logger.Error("Invalid output type", zap.String("type", string(cfg.Output.Type)))
return fmt.Errorf("invalid output type: %s", cfg.Output.Type)
Expand Down
1 change: 1 addition & 0 deletions docker/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,7 @@ docker compose -f docker/docker-compose.telemetry-generator.yml up -d
| `blitz-hostmetrics-linux` | Host Metrics | Synthetic host metrics (CPU, memory, disk, etc.) for Linux |
| `blitz-hostmetrics-windows` | Host Metrics | Synthetic host metrics (CPU, memory, disk, etc.) for Windows |
| `blitz-hostmetrics-promrw` | Host Metrics | Synthetic host metrics sent via the Prometheus remote-write output |
| `blitz-hostmetrics-promscrape` | Host Metrics | Synthetic host metrics served on a Prometheus scrape (pull) endpoint |
| `blitz-traces` | Traces | Synthetic distributed traces (HTTP + DB spans) |
| `blitz-fix` | FIX | FIX protocol messages (4.2 / 4.4 / 5.0 SP2) across 10 asset categories — NewOrderSingle, ExecutionReport, cancel/replace/status |

Expand Down
14 changes: 14 additions & 0 deletions docker/docker-compose.telemetry-generator.yml
Original file line number Diff line number Diff line change
Expand Up @@ -242,6 +242,20 @@ services:
BLITZ_OUTPUT_PROMETHEUS_REMOTE_WRITE_ENDPOINT: ${PROM_RW_ENDPOINT:-http://bdot-collector:19291/api/v1/write}
BLITZ_OUTPUT_PROMETHEUS_REMOTE_WRITE_VERSION: ${PROM_RW_VERSION:-1.0}

# Host metrics served on a Prometheus scrape (pull) endpoint. Point a scraper
# (the collector prometheusreceiver, or a Prometheus server) at
# blitz-hostmetrics-promscrape:9464/metrics to pull the exposition.
blitz-hostmetrics-promscrape:
<<: *blitz-common
environment:
BLITZ_GENERATOR_TYPE: hostmetrics
BLITZ_GENERATOR_HOSTMETRICS_RATE: ${BLITZ_RATE:-1s}
BLITZ_GENERATOR_HOSTMETRICS_OS: linux
BLITZ_OUTPUT_TYPE: prometheus-scrape
BLITZ_OUTPUT_PROMETHEUS_SCRAPE_LISTENADDRESS: 0.0.0.0:9464
ports:
- "${PROM_SCRAPE_PORT:-9464}:9464"

networks:
telemetry-net:
driver: bridge
Expand Down
57 changes: 57 additions & 0 deletions docs/output/prometheus-scrape.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
# Prometheus Scrape Output

The prometheus-scrape output is a metrics-only pull endpoint. It hosts an HTTP `/metrics` endpoint in Prometheus text exposition format, so a Prometheus server or the collector `prometheusreceiver` can scrape blitz directly. Unlike the push-based [remote-write output](prometheus-remote-write.md), the scraper controls timing: blitz holds the current value of each series and serves a snapshot on every scrape.

## Metric Mapping

Generated metric points map to Prometheus series before encoding:

- Gauge and non-monotonic Sum become a gauge.
- Counter becomes a counter with a `_total` suffix.
- Histogram becomes cumulative `_bucket` series with `le` labels, including `+Inf`, plus `_sum` and `_count`.

Metric and label names are sanitized to the Prometheus grammar. Each series carries its name plus the metric's attributes as labels. The endpoint keeps the latest value per series (name, type, and label set), so repeated writes to the same series overwrite rather than accumulate, as a real exporter does. A series not updated within `metricExpiration` is dropped, so a host that stops reporting stops being exposed.

Resource attributes follow the OpenTelemetry-to-Prometheus convention. Every series gets `job` and `instance` labels, taken from `service.namespace`/`service.name` and `service.instance.id`, or from `telemetry.source` and `host.name` when those are absent. A single `target_info` series per target carries the remaining resource attributes.

By default no explicit timestamp is written, and the scraper stamps each sample at scrape time (the idiomatic exporter behavior). Set `emitTimestamps: true` to append each sample's millisecond timestamp instead.

## Configuration

| YAML Path | Flag | Environment Variable | Default | Description |
|---------------------------------------------|-----------------------------------------------|---------------------------------------------------|----------------|---------------------------------------------------------------------------|
| `output.type` | `--output-type` | `BLITZ_OUTPUT_TYPE` | `nop` | Set to `prometheus-scrape` to use this output. |
| `output.prometheus-scrape.listenAddress` | `--output-prometheus-scrape-listenaddress` | `BLITZ_OUTPUT_PROMETHEUS_SCRAPE_LISTENADDRESS` | `0.0.0.0:9464` | Host:port the metrics endpoint binds to. |
| `output.prometheus-scrape.metricsPath` | `--output-prometheus-scrape-metricspath` | `BLITZ_OUTPUT_PROMETHEUS_SCRAPE_METRICSPATH` | `/metrics` | URL path the exposition is served on. |
| `output.prometheus-scrape.emitTimestamps` | `--output-prometheus-scrape-emittimestamps` | `BLITZ_OUTPUT_PROMETHEUS_SCRAPE_EMITTIMESTAMPS` | `false` | Append per-sample millisecond timestamps. |
| `output.prometheus-scrape.metricExpiration` | `--output-prometheus-scrape-metricexpiration` | `BLITZ_OUTPUT_PROMETHEUS_SCRAPE_METRICEXPIRATION` | `5m` | Drop a series not updated within this duration. `0` keeps series forever. |

## Example Configuration

```yaml
output:
type: prometheus-scrape
prometheus-scrape:
listenAddress: 0.0.0.0:9464
metricsPath: /metrics
```

Point a scraper at it, for example a Prometheus scrape config:

```yaml
scrape_configs:
- job_name: blitz
static_configs:
- targets: ["blitz-host:9464"]
```

## Self-Telemetry

blitz records these instruments about the output's own operation and exports them through its telemetry pipeline. They are not the metrics the output exposes to scrapers.

- **`blitz.output.prometheus_scrape.scrapes`** (Counter): scrape requests served.
- **`blitz.output.prometheus_scrape.series_exposed`** (Gauge): series exposed at the last scrape.
- **`blitz.output.prometheus_scrape.exposition_bytes`** (Histogram): exposition body size per scrape.
- **`blitz.output.prometheus_scrape.scrape_latency`** (Histogram, ms): latency of serving a scrape request.

The shared output instruments also record entries received and active workers, tagged `output_type=prometheus-scrape`.
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ require (
github.com/hashicorp/golang-lru/v2 v2.0.7
github.com/jonboulle/clockwork v0.5.0
github.com/prometheus/client_golang v1.24.1
github.com/prometheus/common v0.70.1
github.com/prometheus/prometheus v0.308.1
github.com/spf13/cobra v1.10.2
github.com/spf13/pflag v1.0.10
Expand Down Expand Up @@ -86,7 +87,6 @@ require (
github.com/pb33f/ordered-map/v2 v2.3.1 // indirect
github.com/pelletier/go-toml/v2 v2.2.4 // indirect
github.com/prometheus/client_model v0.6.2 // indirect
github.com/prometheus/common v0.70.1 // indirect
github.com/prometheus/otlptranslator v1.0.0 // indirect
github.com/prometheus/procfs v0.21.1 // indirect
github.com/russross/blackfriday/v2 v2.1.0 // indirect
Expand Down
10 changes: 9 additions & 1 deletion internal/config/output.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,8 @@ const (
OutputTypeHEC OutputType = "hec"
// OutputTypePrometheusRemoteWrite represents Prometheus remote-write output
OutputTypePrometheusRemoteWrite OutputType = "prometheus-remote-write"
// OutputTypePrometheusScrape represents Prometheus scrape (pull) output
OutputTypePrometheusScrape OutputType = "prometheus-scrape"
)

// Output contains configuration for output destinations
Expand All @@ -48,6 +50,8 @@ type Output struct {
Stdout StdoutOutputConfig `yaml:"stdout,omitempty" mapstructure:"stdout,omitempty"`
// PrometheusRemoteWrite contains Prometheus remote-write output configuration
PrometheusRemoteWrite PrometheusRemoteWriteOutputConfig `yaml:"prometheus-remote-write,omitempty" mapstructure:"prometheus-remote-write,omitempty"`
// PrometheusScrape contains Prometheus scrape output configuration
PrometheusScrape PrometheusScrapeOutputConfig `yaml:"prometheus-scrape,omitempty" mapstructure:"prometheus-scrape,omitempty"`
}

// Validate validates the output configuration
Expand Down Expand Up @@ -92,8 +96,12 @@ func (o *Output) Validate() error {
if err := o.PrometheusRemoteWrite.Validate(); err != nil {
return fmt.Errorf("prometheus-remote-write output validation failed: %w", err)
}
case OutputTypePrometheusScrape:
if err := o.PrometheusScrape.Validate(); err != nil {
return fmt.Errorf("prometheus-scrape output validation failed: %w", err)
}
default:
return fmt.Errorf("invalid output type: %s, must be one of: nop, stdout, tcp, udp, syslog, otlp-grpc, file, hec, prometheus-remote-write", o.Type)
return fmt.Errorf("invalid output type: %s, must be one of: nop, stdout, tcp, udp, syslog, otlp-grpc, file, hec, prometheus-remote-write, prometheus-scrape", o.Type)
}

return nil
Expand Down
50 changes: 50 additions & 0 deletions internal/config/output_prometheus_scrape.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
package config

import (
"fmt"
"net"
"strings"
"time"
)

// Default prometheus-scrape output configuration values.
const (
// DefaultPromScrapeListenAddress is the OTel Prometheus exporter default port.
DefaultPromScrapeListenAddress = "0.0.0.0:9464"
// DefaultPromScrapeMetricsPath is the conventional exposition path.
DefaultPromScrapeMetricsPath = "/metrics"
// DefaultPromScrapeMetricExpiration matches the collector prometheusexporter.
DefaultPromScrapeMetricExpiration = 5 * time.Minute
)

// PrometheusScrapeOutputConfig contains configuration for the prometheus-scrape
// output: an HTTP /metrics endpoint served in Prometheus text exposition format.
type PrometheusScrapeOutputConfig struct {
// ListenAddress is the host:port the metrics endpoint binds to.
ListenAddress string `yaml:"listenAddress,omitempty" mapstructure:"listenAddress,omitempty"`
// MetricsPath is the URL path the exposition is served on (e.g. /metrics).
MetricsPath string `yaml:"metricsPath,omitempty" mapstructure:"metricsPath,omitempty"`
// EmitTimestamps, when true, appends each sample's millisecond timestamp to
// the exposition. Default false lets the scraper stamp at scrape time.
EmitTimestamps bool `yaml:"emitTimestamps,omitempty" mapstructure:"emitTimestamps,omitempty"`
// MetricExpiration drops a series not updated within it, so churned or
// stopped hosts stop being exposed. 0 keeps series forever.
MetricExpiration time.Duration `yaml:"metricExpiration,omitempty" mapstructure:"metricExpiration,omitempty"`
}

// Validate validates the prometheus-scrape output configuration. Empty fields
// are allowed: the override system fills defaults before validation.
func (c *PrometheusScrapeOutputConfig) Validate() error {
if c.ListenAddress != "" {
if _, _, err := net.SplitHostPort(c.ListenAddress); err != nil {
return fmt.Errorf("prometheus-scrape output listen address is not a valid host:port: %w", err)
}
}
if c.MetricExpiration < 0 {
return fmt.Errorf("prometheus-scrape output metric expiration cannot be negative, got %s", c.MetricExpiration)
}
if c.MetricsPath != "" && !strings.HasPrefix(c.MetricsPath, "/") {
return fmt.Errorf("prometheus-scrape output metrics path must start with '/', got %q", c.MetricsPath)
}
return nil
}
32 changes: 32 additions & 0 deletions internal/config/output_prometheus_scrape_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
package config

import (
"testing"
"time"
)

func TestPrometheusScrapeOutputConfigValidate(t *testing.T) {
tests := []struct {
name string
cfg PrometheusScrapeOutputConfig
wantErr bool
}{
{"valid", PrometheusScrapeOutputConfig{ListenAddress: "0.0.0.0:9464", MetricsPath: "/metrics"}, false},
{"empty uses defaults", PrometheusScrapeOutputConfig{}, false},
{"emit timestamps ok", PrometheusScrapeOutputConfig{ListenAddress: "127.0.0.1:9464", MetricsPath: "/m", EmitTimestamps: true}, false},
{"path missing leading slash", PrometheusScrapeOutputConfig{ListenAddress: "0.0.0.0:9464", MetricsPath: "metrics"}, true},
{"listen address no port", PrometheusScrapeOutputConfig{ListenAddress: "0.0.0.0", MetricsPath: "/metrics"}, true},
{"negative metric expiration", PrometheusScrapeOutputConfig{MetricExpiration: -time.Second}, true},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
err := tt.cfg.Validate()
if tt.wantErr && err == nil {
t.Fatalf("expected error, got nil")
}
if !tt.wantErr && err != nil {
t.Fatalf("expected no error, got %v", err)
}
})
}
}
4 changes: 4 additions & 0 deletions internal/config/override.go
Original file line number Diff line number Diff line change
Expand Up @@ -359,6 +359,10 @@ func DefaultOverrides() []*Override {
NewOverride("output.prometheus-remote-write.batchSize", "Prometheus remote-write series per batch before flush", DefaultPromRWBatchSize),
NewOverride("output.prometheus-remote-write.batchTimeout", "Prometheus remote-write partial-batch flush timeout", DefaultPromRWBatchTimeout),
NewOverride("output.prometheus-remote-write.timeout", "Prometheus remote-write per-request HTTP timeout", DefaultPromRWTimeout),
NewOverride("output.prometheus-scrape.listenAddress", "Prometheus scrape endpoint listen address (host:port)", DefaultPromScrapeListenAddress),
NewOverride("output.prometheus-scrape.metricsPath", "Prometheus scrape endpoint metrics path", DefaultPromScrapeMetricsPath),
NewOverride("output.prometheus-scrape.emitTimestamps", "append per-sample millisecond timestamps to the exposition", false),
NewOverride("output.prometheus-scrape.metricExpiration", "drop a series not updated within this duration (0 keeps forever)", DefaultPromScrapeMetricExpiration),
NewOverride("telemetry.traces.otlpEndpoint", "OTLP gRPC endpoint (host:port) for exporting blitz's own spans (empty = disabled)", ""),
NewOverride("telemetry.traces.insecure", "send blitz's own spans over plaintext gRPC (no TLS)", false),
NewOverride("telemetry.traces.perBatchSpans", "enable higher-volume per-emit-cycle spans (off by default)", false),
Expand Down
25 changes: 25 additions & 0 deletions internal/config/override_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,10 @@ func getTestOverrideFlagsArgs() []string {
"--output-prometheus-remote-write-batchsize", "250",
"--output-prometheus-remote-write-batchtimeout", "10s",
"--output-prometheus-remote-write-timeout", "20s",
"--output-prometheus-scrape-listenaddress", "127.0.0.1:19464",
"--output-prometheus-scrape-metricspath", "/flagmetrics",
"--output-prometheus-scrape-emittimestamps", "true",
"--output-prometheus-scrape-metricexpiration", "7m",
"--output-stdout-flushinterval", "50ms",
"--metrics-port", "8080",
"--telemetry-traces-otlpendpoint", "traces.example:4317",
Expand Down Expand Up @@ -279,6 +283,10 @@ func getTestOverrideEnvs() map[string]string {
"BLITZ_OUTPUT_PROMETHEUS_REMOTE_WRITE_BATCHSIZE": "300",
"BLITZ_OUTPUT_PROMETHEUS_REMOTE_WRITE_BATCHTIMEOUT": "7s",
"BLITZ_OUTPUT_PROMETHEUS_REMOTE_WRITE_TIMEOUT": "25s",
"BLITZ_OUTPUT_PROMETHEUS_SCRAPE_LISTENADDRESS": "127.0.0.1:29464",
"BLITZ_OUTPUT_PROMETHEUS_SCRAPE_METRICSPATH": "/envmetrics",
"BLITZ_OUTPUT_PROMETHEUS_SCRAPE_EMITTIMESTAMPS": "true",
"BLITZ_OUTPUT_PROMETHEUS_SCRAPE_METRICEXPIRATION": "9m",
"BLITZ_OUTPUT_HEC_ENABLE_TLS": "false",
"BLITZ_OUTPUT_HEC_TLS_CERT": "/env/hec_cert.pem",
"BLITZ_OUTPUT_HEC_TLS_KEY": "/env/hec_key.pem",
Expand Down Expand Up @@ -482,6 +490,11 @@ func TestOverrideDefaults(t *testing.T) {
BatchTimeout: DefaultPromRWBatchTimeout,
Timeout: DefaultPromRWTimeout,
},
PrometheusScrape: PrometheusScrapeOutputConfig{
ListenAddress: DefaultPromScrapeListenAddress,
MetricsPath: DefaultPromScrapeMetricsPath,
MetricExpiration: DefaultPromScrapeMetricExpiration,
},
},
Metrics: Metrics{
Port: DefaultMetricsPort,
Expand Down Expand Up @@ -705,6 +718,12 @@ func TestOverrideFlags(t *testing.T) {
BatchTimeout: 10 * time.Second,
Timeout: 20 * time.Second,
},
PrometheusScrape: PrometheusScrapeOutputConfig{
ListenAddress: "127.0.0.1:19464",
MetricsPath: "/flagmetrics",
EmitTimestamps: true,
MetricExpiration: 7 * time.Minute,
},
},
Metrics: Metrics{
Port: 8080,
Expand Down Expand Up @@ -932,6 +951,12 @@ func TestOverrideEnvs(t *testing.T) {
BatchTimeout: 7 * time.Second,
Timeout: 25 * time.Second,
},
PrometheusScrape: PrometheusScrapeOutputConfig{
ListenAddress: "127.0.0.1:29464",
MetricsPath: "/envmetrics",
EmitTimestamps: true,
MetricExpiration: 9 * time.Minute,
},
},
Metrics: Metrics{
Port: 9100,
Expand Down
Loading
Loading