04QuestionsDistributed counters
04 · Worked prompt
Distributed counters
Design high-write counters with sharding, approximation, aggregation, and monotonic reads.What a passing answer must show
100 points · 45 minutes
- 20pts
Scope the problem
0–5 minPrioritize the core flows, state the scale, and name the non-goals.
- 15pts
Define contracts
5–10 minIdentify durable entities, APIs, idempotency, and the source of truth.
- 30pts
Complete the diagram
10–25 minTrace one write path and one read path. Label the commit boundary and async work.
- 20pts
Lead one deep dive
25–38 minChoose the highest-risk trade-off and explain the mechanism, alternative, and cost.
- 15pts
Prove reliability
38–45 minWalk a failure, recovery, metric, bottleneck, and evolution path.
One complete box-and-arrow design
Design high-write counters with sharding, approximation, aggregation, and monotonic reads.

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 ↗Explain every boundary before adding more boxes.
Writer-sharded deltas absorb updates; watermark-stamped snapshots define read freshness.
Design high-write counters with sharding, approximation, aggregation, and monotonic reads.
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.
End-to-end walkthrough
Trace the architecture in this order.
- 01
Enter and classify the request
Event producers → Counter routerIncrement(op_id, delta) enters over HTTPS / RPC. Counter router handles identity, admission, routing, and request context; it deliberately does not own domain truth.
- 02
Validate, then cross the commit boundary
Counter router → Shard owner → Delta log shardsShard 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.
- 03
Move replayable work off the request path
Shard owner → Aggregation stream → Counter aggregators → Counter snapshotsShard 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.
- 04
Serve reads from the right authority
Counter router → Counter query → Counter snapshots / Delta log shardsCounter 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.
Ownership ledger
Why each box exists—and what it must defend.
| Component | Owns | Why it exists | Interviewer probe |
|---|---|---|---|
| Counter routerCounter + writer hash | Identity, admission, routing | Protects the system edge and attaches trusted context before domain work begins. | Timeout budgets, quotas, regional routing |
| Shard ownerDedupe + append delta | Write invariants and retry identity | Serializes or conditionally applies state changes before acknowledging success. | Concurrent writes, deduplication, hot ownership |
| Delta log shardsDurable increments | Authoritative durable state | Provides the one record used to resolve disputes, recover, and rebuild projections. | Partition key, replication, consistency |
| Aggregation streamShard watermark events | Durable asynchronous handoff | Absorbs bursts and lets slow or optional work retry independently of the request. | Ordering key, lag, retention, dead letters |
| Counter aggregatorsMerge by epoch | Replayable processing | Runs expensive, fan-out, or side-effecting work with leases and bounded retries. | Idempotency, poison work, autoscaling |
| Counter snapshotsTotal + error contract | Rebuildable query state | Shapes data for the dominant reads without weakening the write-side invariant. | Freshness, versioning, rebuild time |
| Counter querySnapshot or fan-out read | Read composition and freshness policy | Chooses authoritative or derived state and returns a stable client contract. | Fan-out, cache policy, partial results |
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.
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 keyAccuracy
Use exact shard sums for billing and sketches for trend discovery
The product must choose the error contractReset
Create a new epoch instead of racing a destructive zero write
Old increments must not leak into the new periodFailure pressure test
Show detection, containment, recovery, and evidence.
Aggregator lag
Serve last watermark and expose freshness
watermark ageDuplicate increment
Deduplicate operation ID for the audit horizon
duplicate deltasShard loss
Replay durable deltas into a replacement shard
recovery duration and unreplayed events- 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
A four-part talk track
- Scope
“I’ll prioritize absorb very high increment rates and read totals with a stated accuracy and freshness.”
- Scale
“The design changes around millions increments/s · skewed keys · seconds freshness · explicit error budget.”
- Decision
“Sharded exact counters trade cheap writes for expensive reads; sketches trade bounded error for fixed memory.”
- 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
Each transition must be durable, observable, and safe to retry.
02Data model and APIs4 entities · 3 interfaces
Core entities
operation_id, counter_id, delta, occurred_atOwner: Write logcounter_id, shard_id, value, versionOwner: Counter storecounter_id, total, watermark, error_boundOwner: Aggregatorcounter_id, epoch, starts_atOwner: Control planeExternal interfaces
/v1/counters/{id}/incrementsApply an idempotent delta
/v1/counters/{id}?consistency=freshRead total, watermark, and error bound
/v1/counters/{id}/resetsBegin 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 keyAccuracy
Use exact shard sums for billing and sketches for trend discovery
The product must choose the error contractReset
Create a new epoch instead of racing a destructive zero write
Old increments must not leak into the new period04Failures, recovery, and evidence3 scenarios
Aggregator lag
Serve last watermark and expose freshness
watermark ageDuplicate increment
Deduplicate operation ID for the audit horizon
duplicate deltasShard loss
Replay durable deltas into a replacement shard
recovery duration and unreplayed events05What 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.
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.