When One Prometheus Isn't Enough: Sharding, Replication, and Thanos — Proven Side by Side
Prometheus doesn’t scale horizontally by itself, and the way it fails is
sneaky. Add a second replica for high availability and every query that
hits both of them sees every series twice — your sum() panels double
overnight. Shard the target list to spread scrape load and every query
that hits one shard sees a slice of the truth and calls it complete —
count(up) says 18 targets when the cluster has 35. Both failure modes
are silent: nothing errors, no pod restarts, the dashboards just quietly
lie.
The canonical answer is Thanos, and every
introduction to it shows you the same architecture diagram and asks you
to take the rest on faith. This lab doesn’t. It runs a plain,
single-server kube-prometheus-stack install and a sharded, replicated,
Thanos-fronted install side by side in one local
kind cluster, then sets a TypeScript
verifier loose on both: same answers as the control group, shards
provably partial, duplicates provably collapsed, and pods killed
mid-run while the queries keep coming back complete. One make up from
a fresh clone builds the whole thing; make verify runs the proofs.
Versions are pinned throughout — kindest/node:v1.36.1 (Kubernetes
1.36), kube-prometheus-stack 88.5.4, Thanos v0.42.4 — so what you
see here keeps working when the charts move on.
The two failure modes, precisely
Replication and sharding break querying for different reasons, and it’s worth being exact about both because Thanos fixes them with two different mechanisms.
Replication runs N identical Prometheus servers scraping the same targets. Each one is a complete, independent copy — great for surviving a crash, useless for querying, because there is no built-in way to ask “the pair” a question. Point Grafana at a load balancer across them and each refresh hits a different server whose scrape timestamps don’t line up; aggregate across both and everything double-counts.
Sharding splits the target list so each server scrapes a subset.
The Prometheus Operator does this natively with a shards field —
targets are hashed across instances — which solves the collection
scaling problem beautifully and creates a query problem: no single
server can answer a cluster-wide question anymore, and none of them
will warn you that its answer is partial.
Thanos addresses both with two components this lab uses (and no more): a sidecar next to each Prometheus pod that exposes that server’s local data over a gRPC StoreAPI, and Thanos Query, a stateless fan-out layer that queries all sidecars, merges the results, and collapses replica duplicates using a designated label. The lab deliberately skips object storage, store gateway, and compactor — that’s the long-term-retention story, a different concept for a different day.
The build: two installs, one cluster
From the lab README, the shape of the whole thing:
kind cluster "lab-thanos" (1 control-plane + 2 workers, kindest/node:v1.36.1)
host ports, mapped straight through by kind:
30900 30901 / 30902 30903 30904
| | | | |
v v v v v
+------------------------+ +-------------------------+ +--------------+ +---------+
| ns: monitoring-vanilla | | ns: monitoring-thanos | | Thanos Query | | Grafana |
| | | | | (2 pods) | | (1 pod) |
| kps-vanilla Prometheus | | shard 0: 2 replicas | +--------------+ +---------+
| (1 pod, no sidecar) | | each + Thanos sidecar |
+------------------------+ | shard 1: 2 replicas |
| each + Thanos sidecar |
+-------------------------+
gRPC StoreAPI (all 4 sidecars) ---->
re-resolved via DNS-SD every 5s
monitoring-vanilla is the control group: one Prometheus, no sharding,
no replication, no Thanos — the setup most teams start with, and the
baseline every proof compares against. monitoring-thanos runs the
same chart version with the same workload split across 2 shards × 2
replicas: four Prometheus pods, each target scraped by exactly one
shard, each shard’s data held twice.
The remarkable part is how small the diff between the two installs is.
Everything that turns a vanilla install into a sharded, replicated,
sidecar-equipped one lives in helm/values-thanos.yaml:
fullnameOverride: kps-thanos
crds:
enabled: false # the vanilla release already installed the CRDs
# ...
prometheus-node-exporter:
service:
port: 9101 # the vanilla install owns 9100 on every host
targetPort: 9101
prometheus:
thanosService:
enabled: true # headless gRPC service across all 4 sidecars
prometheusSpec:
shards: 2 # each scrape target lands on exactly one shard
replicas: 2 # each shard is stored twice
scrapeInterval: 15s
evaluationInterval: 15s
retention: 2h
resources:
requests:
cpu: 100m
memory: 400Mi
thanos:
image: quay.io/thanos/thanos:v0.42.4
Three lines — shards, replicas, thanos.image — are the entire
scaling story on the collection side. The operator handles the rest:
it creates four StatefulSet-managed pods, hashes the scrape targets
across the two shards, injects a Thanos sidecar container into each
pod, and stamps each server’s identity into external labels.
thanosService.enabled: true adds a headless Service that fronts all
four sidecars’ gRPC ports — which matters in a moment.
The surrounding lines are what two parallel installs of the same chart
on one cluster actually cost. The chart’s CRDs are cluster-scoped, so
the second release must not install them again (crds.enabled: false).
node-exporter runs on host networking, so the second DaemonSet has to
move off port 9100 or it would never bind. And both operators run with
prometheusOperator.namespaces.releaseNamespace: true, so each one
only reconciles its own namespace and the installs can’t step on each
other.
Thanos Query is not part of the chart — it’s a plain pinned Deployment
in thanos/query.yaml, and its args are the heart of the whole lab:
containers:
- name: thanos-query
image: quay.io/thanos/thanos:v0.42.4
args:
- query
- --http-address=0.0.0.0:9090
- --grpc-address=0.0.0.0:10901
- --query.replica-label=prometheus_replica
- --endpoint=dnssrv+_grpc._tcp.kps-thanos-thanos-discovery.monitoring-thanos.svc.cluster.local
- --store.sd-dns-interval=5s
Two flags do all the work. The dnssrv+ endpoint points at the
headless service the chart created, so Thanos Query discovers all four
sidecars through DNS SRV records and re-resolves every 5 seconds —
no sidecar addresses hardcoded anywhere, and a killed pod drops out of
the store list within seconds. And --query.replica-label names the
external label that identifies a replica: the operator sets
prometheus_replica=$(POD_NAME) on every server, so the two copies of
each series differ only in that label, and telling Thanos Query which
label is the replica label is exactly what lets it collapse the pair
into one logical series instead of treating them as different data.
One more piece of wiring exists purely to make the failure mode
visible. thanos/shard-services.yaml adds a NodePort service per
shard, selecting on a label the operator puts on each shard’s pods:
apiVersion: v1
kind: Service
metadata:
name: prometheus-shard-0
namespace: monitoring-thanos
spec:
type: NodePort
selector:
operator.prometheus.io/name: kps-thanos-prometheus
operator.prometheus.io/shard: "0"
ports:
- name: http-web
port: 9090
targetPort: 9090
nodePort: 30901
That gives you (and the verifier) a way to interrogate shard 0 and shard 1 directly, bypassing Thanos — which is how you catch a shard confidently returning a wrong answer.
Proof 1: same questions, same answers
Sharding and replicating the collection layer is only acceptable if it
changes nothing about the answers. The verifier’s equivalence check
asks both stacks the same six questions about the cluster and demands
matching results. From verifier/src/commands/equivalence.ts:
// Every query aggregates over the cluster itself (nodes, kubelet cgroups) —
// never over the monitoring pipelines, whose target sets legitimately differ.
export const EQUIVALENCE_CHECKS: EquivalenceCheck[] = [
{ name: "node-count", query: "count(node_uname_info)", kind: "exact" },
{
name: "node-exporter-targets",
query: 'count(up{job="node-exporter"} == 1)',
kind: "exact",
},
{
name: "total-memory",
query: "sum(node_memory_MemTotal_bytes)",
kind: "exact",
},
{
name: "available-memory",
query: "sum(node_memory_MemAvailable_bytes)",
kind: "approx",
},
// ...
];
Counts and fixed totals must match exactly; rate-based numbers get a small tolerance because the two stacks scrape on independent schedules. Making this comparison deterministic across two independent scrape pipelines took one trick worth stealing — both sides are evaluated at the same aligned timestamp, snapped to the shared 15-second scrape grid and set slightly in the past so both stacks have certainly ingested it:
const timeSec = Math.floor(Date.now() / 1000 / 15) * 15 - 30;
Every Thanos-side query in this lab runs with
dedup=true&partial_response=false. That second parameter matters:
it tells Thanos to fail rather than return a partial result when a
store is unreachable, so a passing equivalence check can’t be
Thanos papering over a missing shard — it’s also, quietly, the first
completeness proof.
Proof 2: each shard knows a slice; Thanos sees everything
This is the sharding failure mode made visible on purpose. The
verifier runs count(up) directly against shard 0, directly against
shard 1, and against Thanos Query, then asserts an exact partition —
from verifier/src/commands/partiality.ts:
export function assertPartiality({
shard0,
shard1,
total,
}: {
shard0: number;
shard1: number;
total: number;
}): void {
if (shard0 < 1 || shard1 < 1) {
throw new Error(
`each shard must scrape at least 1 target (shard0=${shard0}, shard1=${shard1})`,
);
}
if (shard0 >= total || shard1 >= total) {
throw new Error(
`each shard must see a strict subset (shard0=${shard0}, shard1=${shard1}, total=${total})`,
);
}
if (shard0 + shard1 !== total) {
throw new Error(
`shard counts must sum to the total (${shard0}+${shard1} != ${total})`,
);
}
}
The assertions are structural rather than hardcoded, because the hash split genuinely varies between runs: one run partitioned 35 targets as 18 + 17, another as 20 + 14. What never varies is the shape — each raw shard returns a confident, non-empty, wrong answer to “how many targets are up?”, and the two wrong answers sum exactly to the right one that Thanos Query returns. If you’ve ever pointed a dashboard at one member of a sharded pair without realizing it, this is the proof of what you were looking at.
Proof 3: six series in, three series out
The replication failure mode and its fix fit in two queries. The
verifier asks Thanos Query for up{job="node-exporter"} twice — once
raw, once deduplicated. From verifier/src/commands/dedup.ts:
const raw = await ctx.thanos.instantQuery('up{job="node-exporter"}', {
dedup: false,
partialResponse: false,
});
const deduped = await ctx.thanos.instantQuery(
'up{job="node-exporter"}',
{ dedup: true, partialResponse: false },
);
With three kind nodes and two replicas per shard, the raw view returns
six series — every node-exporter instance appears once per replica
of the shard that scrapes it, each copy carrying a distinct
prometheus_replica label. The deduplicated view returns exactly
three, one per node. The check doesn’t stop at counting: it groups
the raw series by instance and asserts each one appears under exactly
two distinct replica labels, and that dedup dropped the duplicates
without dropping any instance. Six-to-three, no drops — the
double-counting problem and its fix, side by side in one command.
Proof 4: kill a replica mid-query
The whole point of paying for two replicas per shard is that losing
one is a non-event. The verifier makes that concrete: it deletes
prometheus-kps-thanos-prometheus-1 through the Kubernetes API, then
hammers Thanos Query every couple of seconds for up to 90 seconds
while the pod is down — an instant target count plus a 60-second range
query that must come back gapless. The judging rule is the most
carefully considered code in the lab, from
verifier/src/commands/failover.ts:
// Errors during the outage are tolerated (DNS SD needs a few seconds to drop
// a dead store); what is never tolerated is a *successful* answer with data
// missing — that would be Thanos lying about completeness.
export function assertOutagePolls(polls: OutagePoll[]): void {
const successes = polls.filter((p) => p.ok);
if (successes.length === 0) {
throw new Error(
`no successful Thanos query during the outage window (${polls.length} polls)`,
);
}
const incomplete = successes.filter((p) => !p.complete);
if (incomplete.length > 0) {
throw new Error(
`${incomplete.length} successful poll(s) returned incomplete data during the outage`,
);
}
}
The asymmetry is deliberate. In the few seconds before DNS service
discovery drops the dead store, Thanos may return errors — that’s
partial_response=false doing its job, refusing to guess. An error is
honest. What would be damning is a successful response with data
missing, because that’s the silent lie this whole lab exists to rule
out — and zero of those occur. The same kill-and-hammer sequence then
repeats against one of the two Thanos Query pods itself, proving the
query layer has no single point of failure either.
Implementation left a scar here worth sharing: kubectl delete --wait=false returns before the pod actually stops, and for a brief
moment the dying pod still reports Ready=True. The first live run of
this proof checked “is the pod back?” immediately, saw the not-yet-dead
pod, and exited with zero outage polls — a green checkmark that had
proven nothing. The committed version first waits for the pod to
actually go un-Ready before opening the polling window. If you write
chaos checks of your own, “confirm the outage started” is a step you
skip at your peril.
Watching it live in Grafana
The verifier prints its proofs once; Grafana shows the equivalence
proof continuously. The lab provisions one dashboard — Vanilla vs
Thanos — same questions, same answers — with three paired panels
(nodes seen, busy CPU cores, available memory), each rendered once
against a Vanilla Prometheus datasource and once against Thanos Query. Healthy cluster, twin curves.
One datasource decision hides a trap worth knowing about. From
helm/values-thanos.yaml:
datasources:
defaultDatasourceEnabled: false # the auto-datasource would round-robin raw shards — misleading
The chart’s auto-provisioned datasource points at a Service that selects all four raw Prometheus pods across both shards, so every dashboard query would land on a random pod and return whichever shard’s partial view that pod holds — the exact broken pattern this lab exists to demonstrate, resurrected by a default. Both datasources are provisioned explicitly instead, with Thanos Query as the default.
Run it
make up # kind cluster + both stack installs + Thanos Query + dashboard (~2 min with images cached, longer on a first pull)
make verify # readiness, equivalence, partiality, dedup, failover, Grafana checks (~3 min, kills pods on purpose)
make down # delete the cluster (a few seconds)
Five Prometheus servers plus Thanos on one machine sounds heavy but
isn’t: each Prometheus requests 100m CPU and 400Mi memory, and about
4 GiB of free Docker headroom covers the lot. Every endpoint is a
plain NodePort mapped straight to the host — vanilla Prometheus on
30900, the raw shards on 30901/30902, Thanos Query on 30903,
Grafana on 30904 — so you can rerun any proof by hand with nothing
but a browser and curl.
The lab repository has the full tree: both values files, the Thanos manifests, and the complete verifier. And the deliberate stopping point doubles as the map forward: every server here still retains only two hours locally, so the natural sequel is pointing the sidecars at object storage, with a Store Gateway to query historical blocks and a Compactor to merge and downsample them — the piece that turns “one logical Prometheus” into “one logical Prometheus with unlimited retention.”