Hash every series to three ingesters

The cluster's front door hashes each series' full label set onto a ring and writes every sample to the same three ingesters, succeeding once two confirm, so losing one ingester loses nothing and no popular metric lands on one machine — at the price of every query asking all ingesters.

Previously

Now that remote-write sends every sample to a cluster, the next question is what the cluster does with a sample so no single machine is overloaded or lost.

Scene 10

Hash every series to three ingesters

  1. Watch
  2. Try it
  3. Predict
  4. Capture
instrumentcollectcluster storequeryalertnotifyTIME TO DETECTno dataexampleREMOTE-WRITE BATCH · ONE SERIEShttp_requests_total{route="/pay",status="500",pod="pod-3"}http_requests_total{route="/pay",status="500",pod="pod-3"}shard key: full label set3 copies · 2 must confirm · memory 3×hash ringhash(all labels)distributorvalidate + limitssends each series to 3 ingestersbatch receivedingester-1memory 44%ingester-2memory 47%ingester-3memory 45%ingester-4memory 42%ingester-5memory 46%ingester-6memory 43%query6 / 6ingesters asked per querya query asks all 6 machinesThe batch is checked against the limits, then hashed onto the ring.
What to watch for

What happens to a sample when it enters a scalable metrics cluster? Remote-write hands the cluster a batch, and something has to decide which machine keeps it. Watch one series go through: the front door checks it against the limits, hashes it onto the ring to pick its machines, and sends it to three of them. Two confirmations are enough to call the write done — the third arrives a moment later.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

Distributor.push
check each series against the limits, then hash it
def push(tenant, batch):
for series in batch:
if bad_labels(series) or too_old(series):
reject(series, 400) # never stored
elif over_rate_limit(tenant):
reject(series, 429) # tenant too fast
else:
if shard_by_metric_name:
key = fnv32a(tenant, series.name)
else:
key = fnv32a(tenant, series.labels)
replicate(key, series)
Distributor.replicate
send to every owner, answer once the quorum confirms
def replicate(key, series):
owners = ring.owners(key, n=copies)
quorum = 2 if copies == 3 else 1
acks, fails = 0, 0
for ingester in send_parallel(owners, series):
if ingester.confirmed:
acks += 1
else:
fails += 1 # dead or too slow
if acks >= quorum:
return OK # a third can be late
return ERROR # 500 to remote-write
Querier.select
a read asks whichever ingesters could hold a match
def select(tenant, matchers):
if shard_by_metric_name and has_name(matchers):
key = fnv32a(tenant, name_of(matchers))
targets = ring.owners(key, n=copies)
else:
targets = all_ingesters
samples = []
for ingester in send_parallel(targets, matchers):
samples += ingester.recent_samples(matchers)
return dedupe(samples)

Where this sits in Metrics / Monitoring System

Scene 10 of 18, in the Scale out act — Remote-write, ingesters, dedup, blocks, split queries.. The cluster's front door hashes each series' full label set onto a ring and writes to three ingesters, succeeding once two confirm — no popular metric lands on one machine, but every query asks them all.

Up next. Now that the cluster keeps three safe copies of whatever it receives, the next question is what happens when the scrapers themselves are redundant and send every sample twice.

All 18 scenes in Metrics / Monitoring System · Every curriculum

Built with Arqly
Every scene in Metrics / Monitoring System builds on the one before it.All 18 Metrics / Monitoring System scenes