04QuestionsDistributed counters

04 · Worked prompt

Distributed counters

Design high-write counters with sharding, approximation, aggregation, and monotonic reads.
30 minInterview blueprint
INTERVIEW RUBRIC

What a passing answer must show

100 points · 45 minutes

  1. 20pts

    Scope the problem

    0–5 min

    Prioritize the core flows, state the scale, and name the non-goals.

  2. 15pts

    Define contracts

    5–10 min

    Identify durable entities, APIs, idempotency, and the source of truth.

  3. 30pts

    Complete the diagram

    10–25 min

    Trace one write path and one read path. Label the commit boundary and async work.

  4. 20pts

    Lead one deep dive

    25–38 min

    Choose the highest-risk trade-off and explain the mechanism, alternative, and cost.

  5. 15pts

    Prove reliability

    38–45 min

    Walk a failure, recovery, metric, bottleneck, and evolution path.

DRAW THIS FIRST

One complete box-and-arrow design

Design high-write counters with sharding, approximation, aggregation, and monotonic reads.

Distributed counters · system architecture
Distributed counters system architecture. Writer-sharded deltas absorb updates; watermark-stamped snapshots define read freshness. Request path: Event producers to Counter router to Shard owner to Delta log shards. Asynchronous path: Aggregation stream to Counter aggregators. Read path: Counter query to Counter snapshots.

Write

Event producers enters through Counter router. Shard owner owns validation and commits the durable record to Delta log shards.

Propagate

Aggregation stream separates the committed write from background work. Counter aggregators can retry safely while it builds Counter snapshots.

Read

Counter query serves from Counter snapshots, then checks authoritative state whenever freshness, policy, or correctness requires it.

Say this first: Writer-sharded deltas absorb updates; watermark-stamped snapshots define read freshness.

Open the full whiteboard ↗
DEFEND THE DIAGRAM

Explain every boundary before adding more boxes.

Writer-sharded deltas absorb updates; watermark-stamped snapshots define read freshness.

INTERVIEW CONTRACT

Design high-write counters with sharding, approximation, aggregation, and monotonic reads.

CAPACITY QUESTIONS TO QUANTIFY

Millions increments/s · skewed keys · seconds freshness · explicit error budget. State average and peak load, stored bytes, bandwidth or open connections, and the growth horizon before choosing a partitioning strategy.

01

End-to-end walkthrough

Trace the architecture in this order.

  1. 01
    Enter and classify the request
    Event producers → Counter router

    Increment(op_id, delta) enters over HTTPS / RPC. Counter router handles identity, admission, routing, and request context; it deliberately does not own domain truth.

  2. 02
    Validate, then cross the commit boundary
    Counter router → Shard owner → Delta log shards

    Shard owner receives the command, checks invariants and retry identity, then uses append delta to update Delta log shards. The user-visible mutation is accepted only after this boundary succeeds.

  3. 03
    Move replayable work off the request path
    Shard owner → Aggregation stream → Counter aggregators → Counter snapshots

    Shard owner emits publish after commit; Counter aggregators uses consume and merge epoch to build Counter snapshots. Consumers must tolerate duplicate delivery and stale retries because this path is asynchronous.

  4. 04
    Serve reads from the right authority
    Counter router → Counter query → Counter snapshots / Delta log shards

    Counter query uses read watermark for the common, read-optimized path and strong read when correctness or repair requires authoritative state. The API must state the freshness promise instead of hiding it.

02

Ownership ledger

Why each box exists—and what it must defend.

ComponentOwnsWhy it existsInterviewer probe
Counter routerCounter + writer hashIdentity, admission, routingProtects the system edge and attaches trusted context before domain work begins.Timeout budgets, quotas, regional routing
Shard ownerDedupe + append deltaWrite invariants and retry identitySerializes or conditionally applies state changes before acknowledging success.Concurrent writes, deduplication, hot ownership
Delta log shardsDurable incrementsAuthoritative durable stateProvides the one record used to resolve disputes, recover, and rebuild projections.Partition key, replication, consistency
Aggregation streamShard watermark eventsDurable asynchronous handoffAbsorbs bursts and lets slow or optional work retry independently of the request.Ordering key, lag, retention, dead letters
Counter aggregatorsMerge by epochReplayable processingRuns expensive, fan-out, or side-effecting work with leases and bounded retries.Idempotency, poison work, autoscaling
Counter snapshotsTotal + error contractRebuildable query stateShapes data for the dominant reads without weakening the write-side invariant.Freshness, versioning, rebuild time
Counter querySnapshot or fan-out readRead composition and freshness policyChooses authoritative or derived state and returns a stable client contract.Fan-out, cache policy, partial results
03

Physical design

Name the database, shard key, indexes, and guarantees.

Database + storage
DynamoDB/Cassandra/Scylla or an LSM KV store holds striped counter cells; Kafka is optional for exact replay; Redis serves recent totals.
Partitioning / sharding
Hash counter_id plus stripe_id to spread a hot logical counter. Change stripe count by epoch and merge stripes on read/aggregation.
Indexes
Primary (counter_id, epoch, stripe_id), unique event_id when dedup is required, and aggregates by counter/window.
Replication + consistency
Choose explicitly: linearizable conditional increment for exact low-rate counters, or commutative eventual stripes for huge write rates.
Cache, queue + recovery
Stream/checkpoint events when audit is needed; cache totals with a watermark and error/freshness bound.
Capacity math
Estimate increments/sec for the hottest counter, number of counters, read rate, allowed drift, and stripe merge cost.
Alternative rejected
One atomic row is exact but becomes a hot partition; striping trades read work and bounded staleness for write scale.
04

Deep-dive candidates

Pick one risk and explain the mechanism, alternative, and cost.

Hot keys

Add per-key write shards and periodically merge; keep shard count in metadata

Database partitioning by counter ID alone preserves the hot key
Accuracy

Use exact shard sums for billing and sketches for trend discovery

The product must choose the error contract
Reset

Create a new epoch instead of racing a destructive zero write

Old increments must not leak into the new period
05

Failure pressure test

Show detection, containment, recovery, and evidence.

Aggregator lag

Serve last watermark and expose freshness

watermark age
Duplicate increment

Deduplicate operation ID for the audit horizon

duplicate deltas
Shard loss

Replay durable deltas into a replacement shard

recovery duration and unreplayed events
Before you finish, explicitly cover
  • Functional requirements and non-goals
  • Peak traffic, storage, bandwidth, and growth
  • Entities, APIs, idempotency, and pagination
  • Source of truth and consistency promise
  • Partition key, replicas, caches, and hot spots
  • Retries, backpressure, failover, and reconciliation
  • Latency, saturation, correctness, and recovery metrics
  • Security, migration, cost, and multi-region evolution
SAY THIS WHILE YOU DRAW

A four-part talk track

  1. Scope

    “I’ll prioritize absorb very high increment rates and read totals with a stated accuracy and freshness.”

  2. Scale

    “The design changes around millions increments/s · skewed keys · seconds freshness · explicit error budget.”

  3. Decision

    “Sharded exact counters trade cheap writes for expensive reads; sketches trade bounded error for fixed memory.”

  4. Risk

    “The first failure I want to pressure-test is: Duplicate delivery and non-atomic reset create permanent drift.”

Reference details

Open these only after you can explain the diagram above without reading.

01Requirements and state lifecycle4 requirements
  • Absorb very high increment rates
  • Read totals with a stated accuracy and freshness
  • Handle hot keys and node loss
  • Support reset, correction, and audit when required
Distributed counters · state lifecycle
02Data model and APIs4 entities · 3 interfaces

Core entities

Incrementoperation_id, counter_id, delta, occurred_atOwner: Write log
CounterShardcounter_id, shard_id, value, versionOwner: Counter store
CounterSnapshotcounter_id, total, watermark, error_boundOwner: Aggregator
ResetEpochcounter_id, epoch, starts_atOwner: Control plane

External interfaces

POST /v1/counters/{id}/increments

Apply an idempotent delta

GET /v1/counters/{id}?consistency=fresh

Read total, watermark, and error bound

POST /v1/counters/{id}/resets

Begin a new versioned epoch

03Deep dives and trade-offsChoose one

Hot keys

Add per-key write shards and periodically merge; keep shard count in metadata

Database partitioning by counter ID alone preserves the hot key

Accuracy

Use exact shard sums for billing and sketches for trend discovery

The product must choose the error contract

Reset

Create a new epoch instead of racing a destructive zero write

Old increments must not leak into the new period
04Failures, recovery, and evidence3 scenarios

Aggregator lag

Serve last watermark and expose freshness

watermark age

Duplicate increment

Deduplicate operation ID for the audit horizon

duplicate deltas

Shard loss

Replay durable deltas into a replacement shard

recovery duration and unreplayed events
05What makes the answer seniorInterviewer signals
  • Ask exact, monotonic, or approximate before designing
  • Return a watermark, not an unexplained number
  • Resets are a versioning problem, not a SET value=0 call
  • Primary trade-off: Sharded exact counters trade cheap writes for expensive reads; sketches trade bounded error for fixed memory.
BEFORE THE NEXT QUESTION

Can you redraw it from memory?

  • Name the source of truth.
  • Trace the write and read paths.
  • Defend one trade-off.
  • Recover from one failure.