Skip to main content

Data Pipeline

ServiceRadar uses two primary data paths:

  • Agent ingestion (edge status, discovery, sync) via mTLS gRPC to agent-gateway
  • Bulk telemetry ingestion (logs/flows/etc.) via NATS JetStream streams into CNPG

This page is intentionally high-level. It focuses on the mental model and the moving parts you will see when debugging.

NATS JetStream (Bulk Ingestion)​

Bulk collectors publish to NATS JetStream. The most common stream is events.

serviceradar_core is the bulk-ingestion writer for these streams. It runs the Zen decision rules in-process for log normalization and persists records into CNPG. EventWriter's logs consumer promotes logs that match an event rule into OCSF events as it stores them; there is no separate promotion consumer, so each log is promoted once, under an id derived from the JetStream message.

For sysmon, SNMP, ICMP, rperf, MTR scalar and sweep metrics, agent-gateway splits oversized protobuf batches before publishing them to NATS. The gateway uses the broker-advertised max_payload with headroom for NATS headers; no agent-side batch-size setting is required. Large multi-port sweep groups can therefore produce several messages on the same subject. A warning reports the original byte count, broker limit, number of parts and dropped_points. A point that cannot fit by itself is dropped; if no parts remain, publishing returns an error. The payload-budget and per-part deduplication contracts live in ServiceRadarAgentGateway.MetricBatchPublisher and ServiceRadarAgentGateway.MetricBatchSplit.

Data Service (datasvc)​

datasvc is a gRPC service (port 50057) that fronts the platform's NATS-backed key-value and object stores. Other components use it for shared configuration and state—for example, the KV bucket serviceradar-datasvc holds rule definitions and runtime settings, and the object store carries larger payloads. Routing this state through one service keeps NATS KV/object access consistent and avoids components manipulating JetStream buckets directly.

CNPG (System Of Record)​

CNPG is the system of record for inventory, telemetry, and analytics (Timescale hypertables and AGE graph features are enabled in the cluster).

An optional StarRocks warehouse (off by default) can hold migrated telemetry history -- flows, scalar metrics, logs, event history, and further append-only datasets, including BMP routing events -- written by the same EventWriter path. The BMP table and retention contract is in k8s/starrocks/README.md. CNPG remains the system of record for inventory, configuration, credentials and current alert state, and keeps serving scalar metrics, logs and event history until each is explicitly cut over. Flows are the exception: they are written to CNPG but served only from the warehouse, so in:flows and the NetFlow dashboard are refused until flows is cut over. See NetFlow for the flow path and Helm Deployment and Configuration for the chart values and where the warehouse install is documented.

Querying happens through the web UI and SRQL, which is embedded in web-ng.