System Design Problem

Design Top K Most Shared Articles

Commonly Asked By:TwitterMetaLinkedInReddit

Interview Setup

Interview Prompt

Design a system that tracks article shares and returns the top-K most shared articles per time window (1h, 24h, 7d, 30d), category, and region. Handle 1B share events/day with rankings updating within 1 to 5 minutes.

Clarifying Questions (ask before designing)

QuestionWhy it matters
Is an exact share count required, or is a share-count estimate within 5% acceptable?Enables Space-Saving and Lossy Counting algorithms instead of maintaining discrete counters for every article across billions of share events.
How do we canonicalize URLs, including protocol differences and tracking query parameters, to ensure accurate counts?Without normalization, shares for a single viral article split across multiple independent counter keys.
Do we need trending velocity detection or just cumulative share volume?Velocity requires sliding windows and z-score calculations, whereas raw count metrics favor older evergreen articles.
What is the anticipated read QPS on top-K endpoints relative to the share write rate?An ingest volume of 1B shares per day (~100K peak/sec) represents a write heavy profile, which requires precomputed top-K views rather than on demand sorting.

Scope

In scope

  • Lossy Counting
  • Time-window leaderboards
  • MapReduce aggregation
  • Approximation algorithms
  • Capacity estimation with shown math

Out of scope (state explicitly)

  • Detailed frontend/UI pixel implementation
  • Org structure, staffing, and hiring plan

Functional Requirements

This design shares core streaming aggregation patterns with Top K Rankings System, adapted specifically for tracking article shares across pane based windows and configurable leaderboard contexts. During requirements gathering, clarify the time windows, filtering dimensions by category and region, and whether share-count estimates within a 5% target are acceptable.

In an interview, cross-reference the Top K Rankings System problem because both architectures rely on the min-heap streaming pattern while handling different payload types.

  • Track shares: Record each article share event with URL, user, timestamp, category, and regional metadata.
  • Top-K ranking: Return the K most shared articles across configured time windows including 1h, 24h, 7d, and 30d.
  • Near real-time freshness: Update leaderboard rankings within 1 to 5 minutes of new share events.
  • Category filtering: Query top-K rankings partitioned by vertical categories such as technology, sports, or politics.
  • Regional filtering: Query top-K rankings partitioned by geographic country or region.
  • Trending detection: Identify and flag emerging articles exhibiting rapid spikes in share velocity.

Non-Functional Requirements

The system must support sub-50ms read latency on precomputed leaderboard views while maintaining ranking freshness within minutes. Because write heavy share workloads peak at scale, ingest processing relies on stream-based aggregation.

  • High Throughput: Sustain the 100K share events per second peak design target, with headroom validated through load testing.
  • Low Latency: Serve top-K read queries in under 50 milliseconds at p99.
  • Accuracy SLA: Share-count estimates within a 5% target are acceptable, which enables bounded memory streaming sketches.
  • Scalability: Support 1B share events per day across tens of millions of distinct articles.
  • High Availability: Target 99.99% uptime for both share ingestion and leaderboard retrieval paths.

Capacity Estimations

Share events per day and the volume of independent leaderboard views drive Flink partitioning and Redis memory sizing. At billions of shares, exact counters for every active article create large O(N) state, which makes it essential to justify bounded approximation algorithms like Space-Saving or Count-Min Sketch early in the conversation.

MetricCalculationValue
Share events / dayGiven1B
Share events / sec1B ÷ 86400~12K (peak 100K)
Unique articles shared / dayGiven10M
Top-K returnedK100 typically
Share event sizeGiven logical wire-event planning assumption100 bytes
Peak events per Kafka partition100K ÷ 64≈ 1,563 events/sec/partition
Time windowsGiven (assumption documented in value)1h, 24h, 7d, 30d
5 minute panes for 7d7 x 24 x 122,016 panes
5 minute panes for 30d30 x 24 x 128,640 panes
7d pane memory example100K active article keys x 2,016 panes x 4 bytes≈ 806 MB before framework overhead
30d pane memory example100K active article keys x 8,640 panes x 4 bytes≈ 3.46 GB before framework overhead
Leaderboard contextswindows x active regions x active categoriesWorkload dependent

Architecture Diagram

The ingestion architecture records incoming share events, applies deterministic URL normalization, aggregates 5 minute pane counts per article, region, and category, and derives 1h, 24h, 7d, and 30d leaderboards through stream processing. Flink maintains bounded candidate summaries per 5 minute pane and leaderboard context, merges candidates from the relevant recent panes, and refreshes Redis every 60 seconds. Redis lookup targets are under 5 milliseconds, while the Top-K API targets under 50 milliseconds at p99. ClickHouse remains off the critical path and stores the durable event history used for exact reconciliation and ad hoc analytics after viral spikes, canonicalization changes, or approximation drift.

In an interview, emphasize that deterministic URL normalization runs before stream aggregation, while shortener resolution can complete asynchronously and reconcile provisional article IDs so shares for the same article do not remain split across multiple keys.

Loading...

Component Deep Dives

Multi-Window Counting with Panes

Pane based counting turns the high volume share stream into fixed 5 minute summaries that can be reused across the 1 hour, 24 hour, 7 day, and 30 day leaderboard windows. Because 100K share events per second cannot trigger global sorts on demand, the platform maintains pane state and derives bounded candidate sets before publishing rankings to Redis.

Fixed tumbling windows suffer from artificial cliff effects at window boundaries. In contrast, pane based counting divides time into discrete 5 minute buckets, enabling 1 hour, 24 hour, 7 day, and 30 day leaderboards to be computed as running sums of recent panes. To conserve memory, only articles that register meaningful velocity enter the active candidate set. This threshold is a heuristic and must be validated against the 5% accuracy target because it can exclude low velocity articles whose shares accumulate across panes.

YAML
pane_based_counting:
  time_slice_unit: "5 minute discrete panes"
  pane_storage: "Aggregated count of shares per article, region, and category within each 5 minute slice"
  window_aggregations:
    one_hour: "Sum of the most recent 12 panes"
    twenty_four_hours: "Sum of the most recent 288 panes"
    seven_days: "Sum of the most recent 2,016 panes"
    thirty_days: "Sum of the most recent 8,640 panes"
  memory_optimization:
    threshold_filter: "Heuristic: only articles with more than 10 shares in any pane enter candidate state. Validate this threshold against the 5% accuracy SLA because it can exclude low-velocity articles whose shares accumulate across panes"
    memory_calculation: "100K active article keys * 2,016 panes * 4 bytes ≈ 806 MB before framework overhead"
    thirty_day_memory_calculation: "100K active article keys * 8,640 panes * 4 bytes ≈ 3.46 GB before framework overhead"
    counter_width_note: "The 4 byte calculation is an illustrative lower-bound estimate, so use a wider counter type when counts can exceed 2^32"

Approximate Top-K: Space-Saving Algorithm

Tracking exact per-article counters across 10M unique URLs per day creates O(N) state, and the memory multiplies across regions, categories, windows, and retention horizons. The Space-Saving algorithm addresses this by maintaining a fixed candidate set of 1,000 entries for each active 5 minute pane and leaderboard context while bounding frequency overestimation. The reference implementation scans 1,000 entries and therefore runs in O(K) per update. A production min-heap implementation can reduce update work to O(log K). When a pane closes, its bounded candidate summary feeds the window aggregator, which merges candidates from the recent panes for the 1h, 24h, 7d, and 30d views. We combine this approximation with hourly ClickHouse reconciliation for the published top-100 leaderboard. The 1,000 candidate capacity is an illustrative starting point. The 5% accuracy target is a product and measurement SLA, not a universal Space-Saving guarantee, so the pipeline must validate observed error against the target and increase candidate capacity or refine candidates when needed.

TYPESCRIPT
// Illustrative fixed-capacity Space-Saving tracker for one leaderboard context.
// Production Flink state should store the sketch in managed state so checkpoints can recover it.
class SpaceSavingTopK {
  private readonly capacity = 1000;
  private entries: Map<string, number> = new Map();

  onShareEvent(articleId: string): void {
    if (this.entries.has(articleId)) {
      // Case 1: Candidate already tracked, increment counter
      this.entries.set(articleId, (this.entries.get(articleId) ?? 0) + 1);
    } else if (this.entries.size < this.capacity) {
      // Case 2: Available capacity, insert candidate with initial count of 1
      this.entries.set(articleId, 1);
    } else {
      // Case 3: Saturated capacity, evict minimum entry and increment its score
      const [minArticleId, minCount] = this.findMinimumEntry();
      this.entries.delete(minArticleId);
      this.entries.set(articleId, minCount + 1);
    }
  }

  // The sketch bounds frequency overestimation, but boundary top-K membership is approximate.
  // This reference implementation scans K entries, so each update is O(K).
  // A production min-heap can reduce minimum lookup/update work to O(log K).
  // Reset the sketch when the 5 minute pane closes, as window state is built from recent pane summaries.
  private findMinimumEntry(): [string, number] {
    let minKey = "";
    let minVal = Number.POSITIVE_INFINITY;
    for (const [key, val] of this.entries.entries()) {
      if (val < minVal) {
        minVal = val;
        minKey = key;
      }
    }
    return [minKey, minVal];
  }
}

Viral Article: Single Key Hot Spot

A viral article can become an early hotspot when 1M shares arrive within 10 minutes and concentrate counter updates on one Redis key and one Flink keyed task. The scenario averages about 1.7K shares per second for that article, so the operational threshold is a tuning trigger for burstiness rather than a claim that the average rate itself exceeds the threshold. The gateway can coalesce serving-counter updates for 1 second while raw share events remain durable in Kafka. Sustained hotspots can instead shard the counter across 8 deterministic sub-keys using hash(event_id) modulo 8, then sum the partial counts in a second aggregation stage. Monitoring tracks per-key QPS and can trigger sharding when an article exceeds an operational threshold such as 10K operations per minute. The threshold is a tuning heuristic rather than a universal limit.

YAML
hotspot_scenario:
  trigger: "Breaking news event generates 1M shares in 10 minutes"
  failure_mode: "All write increments hit a single Flink keyed task and Redis counter shard, causing worker backpressure"

remediation_strategies:
  gateway_microbatch:
    description: "API servers coalesce serving-counter updates in local memory for 1 second"
    batch_flush:
      article_id: "art-123"
      count: 342
    durability_note: "Raw share events continue to Kafka, ensuring gateway memory batching is never the only durable copy"
    result: "Reduces Redis counter update QPS for bursty hot articles"

  sub_key_sharding:
    description: "Partition counter updates using hash(event_id) % 8 across 8 deterministic sub-keys"
    stage_two_merge: "Downstream Flink aggregation sums partial counts across all 8 sub-keys"
    ordering_note: "Global per-article order is relaxed in sharded mode because share counting is commutative, while identical event retries stay on the same shard"
    result: "Distributes a single viral article's counter load across multiple cluster nodes"

Event Bus Design (Kafka)

The Kafka event bus decouples asynchronous ingestion from downstream processing consumers, including window counters, top-K sketch aggregators, and persistent ClickHouse writers. Partitioning topics by article_id preserves in-order delivery for each article in the baseline unsalted path, while checkpointed Flink state preserves the required recovery point across worker restarts. Hot-key sharding intentionally relaxes global per-article ordering because counting is commutative.

YAML
topic_configuration:
  name: "share-events"
  partitions: 64
  partition_key: "article_id"
  ordering: "Preserves per-article order in the baseline unsalted path"
  hot_key_mode: "For sharded viral articles, use article_id + shard_id and rely on commutative counting rather than global ordering"
  retention_days: 7
  replication_factor: 3
  min_insync_replicas: 2

producer:
  source: "Share API after deterministic URL normalization, where unresolved shorteners may use a provisional article_id"
  idempotency: "Derive one event_id per authenticated user plus idempotency key, so retries return the same event_id within the retention window"
  idempotency_retention: "Keep the request idempotency record for at least the maximum expected client retry window"
  acks: "all"
  enable_idempotence: true
  event_schema:
    event_id: "string"
    article_id: "string"
    canonical_url: "string"
    user_id: "string"
    platform: "string"
    region: "string"
    category: "string"
    timestamp: "int64"

consumer_groups:
  window_counter:
    engine: "Apache Flink"
    responsibility: "Maintains 5 minute pane counts supporting 1h, 24h, 7d, and 30d windows per article, region, and category"
  topk_aggregator:
    engine: "Apache Flink"
    responsibility: "Maintains bounded candidate summaries for active 5 minute panes, merges candidates from the relevant panes for each leaderboard context, and flushes ranked results to Redis every 60 seconds"
  persistence_writer:
    engine: "ClickHouse Ingest Worker"
    responsibility: "Uses durable idempotency state keyed by event_id, retries safely, and executes micro-batch INSERT every 10 seconds into ClickHouse for exact historical reconciliation and auditing"

operational_paths:
  read_path: "GET /api/v1/top-articles -> Queries Redis ZREVRANGE topk:24h:US:tech 0 99 WITHSCORES"
  dead_letter_queue:
    topic: "share-events-dlq"
    reasons: "Poison events after bounded retries, schema failures, or unrecoverable processing errors"
  consumer_lag_alert:
    threshold: "60 seconds"

API Design

Ranking API Domain Signatures

The ranking API defines the domain contracts used to return precomputed top-K views from Redis. The read path does not sort the full article catalog on demand, which supports the under 50 millisecond p99 target. The share ingest path uses an idempotency key so client retries do not create duplicate logical share events.

TypeScript domain models define the interface contracts for authenticated share events, time window aggregations, and leaderboard response structures.

TYPESCRIPT
// Core domain contracts for article sharing and top-K leaderboard queries

export type ArticleId = string;
export type ShareEventId = string;
export type UserId = string;
export type CanonicalUrlHash = string;
export type RegionCode = string;
export type Category = string;
export type UnixMs = number;
export type TimeWindow = "1h" | "24h" | "7d" | "30d";
export type SharePlatform = "web" | "ios" | "android";

export interface ShareEventPayload {
  eventId: ShareEventId;
  articleId: ArticleId;
  canonicalUrl: string;
  canonicalUrlHash: CanonicalUrlHash;
  userId: UserId;
  platform: SharePlatform;
  region: RegionCode;
  category: Category;
  sharedAt: UnixMs; // Server-assigned event timestamp in Unix epoch milliseconds
}

export interface IngestShareRequest {
  url: string;
  platform: SharePlatform;
  region: RegionCode;
  category: Category;
  idempotencyKey: string; // Derived from the Idempotency-Key request header
  // userId is derived from the authenticated identity, not accepted from the client.
  // region and category are validated against server-managed allowlists or taxonomy.
}

export interface IngestShareResponse {
  eventId: ShareEventId;
  articleId: ArticleId;
  status: "queued" | "processed";
}

export interface TopArticlesQuery {
  window: TimeWindow;
  region?: RegionCode;
  category?: Category;
  limit?: number; // Default 50, maximum 100
}

export interface RankedArticleEntry {
  articleId: ArticleId;
  url: string;
  title: string;
  shareCount: number;
  rank: number;
}

export interface TopArticlesResponse {
  articles: RankedArticleEntry[];
  window: TimeWindow;
  region?: RegionCode;
  category?: Category;
  approximate: boolean; // True when the published ranking uses a bounded streaming sketch
  computedAt: string; // ISO 8601 timestamp
}

Get Top-K Shared Articles

Retrieves the top-K most shared articles for a specified time window, geographic region, and content category from precomputed Redis sorted sets.

HTTP
GET /api/v1/top-articles?window=24h&region=US&category=tech&limit=50 HTTP/1.1
Host: api.example.com
Authorization: Bearer <access-token>
Accept: application/json

HTTP/1.1 200 OK
Content-Type: application/json

{
  "articles": [
    {
      "article_id": "art-123",
      "url": "https://example.com/news/tech-breakthrough",
      "title": "Major Tech Breakthrough Announced",
      "share_count": 152340,
      "rank": 1
    }
  ],
  "window": "24h",
  "region": "US",
  "category": "tech",
  "approximate": true,
  "computed_at": "2026-03-14T10:05:00Z"
}

Record Article Share

Ingests an article share event asynchronously into the streaming pipeline, validating the payload and returning an HTTP 202 acknowledgment.

HTTP
POST /api/v1/shares HTTP/1.1
Host: api.example.com
Authorization: Bearer <access-token>
Content-Type: application/json
Idempotency-Key: 01HR8Q7M9V4Y6D2P8K3N5T7W9X

{
  "url": "https://example.com/news/tech-breakthrough?utm_source=twitter",
  "platform": "web",
  "region": "US",
  "category": "tech"
}

HTTP/1.1 202 Accepted
Content-Type: application/json

{
  "event_id": "evt-98213",
  "article_id": "art-123",
  "status": "queued"
}

Common Error Responses

Standard error codes returned across the ingestion and leaderboard read endpoints.

400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff
504 Gateway Timeout: search index shard responded slowly, narrow query parameters or retry

Data Model

Redis Storage Schema

The Redis schema supports the active leaderboard view by storing sorted sets for ranking and hashes for article metadata. Durable share history remains in ClickHouse for reconciliation and audit.

REDIS
# Active Leaderboard Sorted Set
leaderboard_sorted_set:
  key_pattern: "topk:{window}:{region}:{category}"
  data_type: "Sorted Set (ZSET)"
  score: "share_count plus fractional tie-breaker less than 1"
  member: "article_id"
  refresh: "Rebuild or trim the published member set every 60 seconds so expired window members are removed"
  ttl: "2x the window duration as fallback key cleanup, because TTL alone does not expire individual ZSET members"
  query_command: "ZREVRANGE topk:24h:US:tech 0 99 WITHSCORES"
  update_command: "ZADD topk:24h:US:tech 152340.000321 art-123"
  read_note: "The integer share count is recovered from the score's integer portion, while the fractional portion is reserved for tie-breaking and remains below 1"

# Article Metadata Cache
metadata_hash:
  key_pattern: "article_meta:{article_id}"
  data_type: "Hash"
  fields:
    url: "https://example.com/news/tech-breakthrough"
    title: "Major Tech Breakthrough Announced"
    image_url: "https://cdn.example.com/images/thumb-123.jpg"
    publisher: "Tech Wire"
    category: "tech"
  ttl: "7 days sliding expiration, refreshed explicitly on successful read or write"

ClickHouse Columnar Storage

ClickHouse provides append only columnar storage for all raw share events, coupled with a SummingMergeTree materialized view that aggregates hourly share counts per article, region, and category for audit and reconciliation.

SQL
CREATE TABLE share_events (
    event_id         String,
    article_id       String,
    canonical_url    String,
    user_id          String,
    platform         LowCardinality(String),
    region           LowCardinality(String),
    category         LowCardinality(String),
    shared_at        DateTime64(3)
) ENGINE = MergeTree()
PARTITION BY toYYYYMMDD(shared_at)
ORDER BY (article_id, region, category, shared_at);

-- The ClickHouse ingest worker deduplicates by stable event_id before insert.
-- Retries of the same client idempotent request or Kafka delivery therefore do not inflate historical counts.
-- ClickHouse is not relying on a uniqueness constraint, because idempotency is enforced in the ingest path.

CREATE MATERIALIZED VIEW share_counts_hourly
ENGINE = SummingMergeTree()
ORDER BY (article_id, region, category, hour)
AS SELECT
    article_id, region, category,
    toStartOfHour(shared_at) AS hour,
    count() AS share_count
FROM share_events
GROUP BY article_id, region, category, hour;

Fault Tolerance

ConcernSolution
Flink failureCheckpoint Flink pane and operator state to S3 every 60 seconds and resume from the latest consistent checkpoint upon worker restart.
Redis lossFlink restores managed pane state from the latest consistent checkpoint and rebuilds the affected top-K Redis sorted sets. The 60-second rebuild target is a scenario SLO rather than a Redis guarantee and excludes checkpoint restoration time when the latest checkpoint is older.
Late eventsFlink bounded-out-of-orderness watermarks advance event time using a 5 minute out-of-order tolerance. Configure allowed lateness and side-output handling for records that arrive after the watermark so late events can be reconciled and corrective leaderboard refreshes can be triggered without silently changing the published approximate count.
Article metadata missingEnqueue asynchronous metadata enrichment workers and fall back to displaying the canonical URL if the title is not yet cached.
Spam sharesApply sliding window rate limits of 100 shares per hour per user account and trigger bot anomaly detection algorithms.

Additional Considerations

URL Deduplication for Articles

Raw article URLs often contain varying protocols, domain aliases, redirect shorteners, and campaign tracking parameters. A canonical normalization pipeline ensures that all shares of a given story resolve to a single deterministic identifier.

YAML
canonical_url_normalization_pipeline:
  step_1_deterministic_normalization: "Normalize scheme and host casing, strip default ports, normalize safe path encoding, and remove only configured tracking parameters such as utm_* and fbclid. Do not assume http and https are equivalent without policy, and do not strip arbitrary query parameters because they may be part of content identity"
  step_2_www_policy: "Do not remove 'www.' globally, treating host aliases as equivalent only when an explicit publisher or domain policy confirms equivalence"
  step_3_hash_generation: "Compute canonical_url_hash = SHA256(deterministically_normalized_url)"
  step_4_identifier_assignment: "Use canonical_url_hash as the initial article_id for stream aggregation and caching keys, preferring a publisher supplied stable article ID when one exists"
  step_5_shortener_resolution: "Resolve HTTP 301, 302, 307, and 308 shorteners asynchronously through an isolated fetcher with timeouts, hop limits, response-size limits, protocol allowlists, and private-network SSRF protections"
  step_6_alias_reconciliation: "When async resolution proves that multiple provisional IDs map to the same destination, merge counts through an alias table and ClickHouse reconciliation job"

Clickbait and Spam Filtering

Viral share counts can be manipulated through bot coordination or syndicated spam networks. The system preserves raw share events for audit while applying suppression, quarantine, or ranking penalties to the serving score rather than rewriting the authoritative historical count.

YAML
spam_detection_rules:
  account_diversity_anomaly:
    condition: "More than 80% of shares originate from fewer than 50 unique user accounts"
    classification: "Potential coordinated bot ring manipulation"
    action: "Suppress the article from the serving leaderboard while retaining raw share events for audit"

  domain_reputation_penalty:
    condition: "Article URL matches known content farm or low-reputation domain catalog"
    action: "Apply 0.1x ranking score multiplier and do not alter the raw historical share count"

  velocity_step_function:
    condition: "Share velocity spikes instantly from 0 to 100K shares within minutes without organic acceleration"
    action: "Quarantine ranking contribution for asynchronous anomaly verification"

  account_age_heuristic:
    condition: "More than 50% of sharing users registered their accounts within the past 7 days"
    action: "Down-rank or quarantine unverified shares as a ranking heuristic and retain raw events for reconciliation"

  policy_note:
    meaning: "These rules are ranking heuristics rather than proof of abuse, so keep raw events for later audit and reprocessing"

Virality Detection

Sudden surges in article sharing indicate breaking news or viral trends that warrant priority indexing, CDN cache pre-warming, and editorial notifications.

YAML
virality_sliding_windows:
  metric: "Shares per minute evaluated continuously over a 5 minute sliding window"
  classification_stages:
    stage_1_emerging:
      threshold: "velocity > 50 shares/min"
      action: "Flag article in internal editorial dashboard"
    stage_2_viral:
      threshold: "velocity > 500 shares/min"
      action: "Publish Kafka event to initiate CDN edge cache pre-warming"
    stage_3_mega_viral:
      threshold: "velocity > 5000 shares/min"
      action: "Dispatch push notifications and surface on home feed carousel"

Related Problems and Concepts

These links provide complementary coverage for leaderboard rankings, event streaming, and velocity detection.

  • Trending Topics: Optimizes for share velocity and sudden bursts using z-score computations on hashtags, whereas this design computes cumulative share volumes across tumbling and sliding windows.
  • Top K Rankings System: Explores global e-commerce and app store bestseller rankings using parallelized min-heap merge operators.
  • Stream Processing Basics: Details tumbling, sliding, and session window semantics, event time processing, and watermark mechanics in Apache Flink.
  • Redis Patterns for Interview Systems: Explains sorted sets, fast ZREVRANGE retrieval, and atomic Lua scripting for high throughput counting.
  • Caching Patterns and Invalidation: Covers TTL policies and cache stampede mitigations when serving hot leaderboard queries at scale.

Interview Walkthrough

Structure the 45 minute discussion around the path from a share click to the precomputed Redis sorted set.

  • 25 minute cut

    Skip secondary operational deep dives unless the interviewer pushes for staff-level depth.

    • Normalize URLs before counting by removing known tracking parameters and handling short link resolution safely (5 min).
    • Establish sliding time windows over the last 24 hours with decay to prioritize active engagement (6 min).
    • Implement the Space-Saving sketch or Count-Min Sketch for bounded approximate streaming counts (5 min).
    • Deploy spam detection heuristics on share velocity anomalies to filter out coordinated bots (5 min).
    • Penalize and quarantine articles where more than 80% of shares originate from low-diversity accounts (4 min).
  • Normalize URLs before counting by stripping tracking parameters, expanding short links, and converting hostnames to lowercase because otherwise shares for the same article split across multiple identifiers.
  • Use pane based windows for 24 hour and longer rankings, and apply time decay only where the product requirement calls for recency weighting rather than lifetime totals.
  • Apply Space-Saving or Count-Min Sketch for approximate per-article counts. The reference Space-Saving implementation is O(K) per event because it scans the candidate set, while a heap-backed implementation can approach O(log K). This avoids an O(N log N) global sort on every share event.
  • Run spam detection on share velocity anomalies because organic growth builds gradually, whereas malicious bot campaigns jump abruptly from 0 to 100K shares within minutes.
  • Penalize articles when more than 80% of shares originate from fewer than 50 unique accounts or from accounts that were created less than 7 days ago.
  • Publish virality stage transitions through Kafka events to trigger automated CDN edge cache prewarming and editorial newsroom alerts.
  • Quantify ingest scale: 100K shares/sec multiplied by 100 bytes per event yields approximately 10 MB/s of logical event payload throughput, before Kafka protocol overhead, replication, and storage overhead. With 64 partitions, that is approximately 1,563 events/sec per partition on average, but partition sizing must still account for uneven article popularity and hot keys.
  • Highlight the common architectural pitfall of executing an exact global sort on every share event, which causes memory and CPU consumption to explode once the catalog exceeds 10M entries.

Engineering Trade-offs

Exact vs Approximate Top-K

Selecting between full cardinality counters and bounded streaming sketches determines the memory cost and accuracy tradeoff for the top-K leaderboard.

ApproachMemoryAccuracySpeed
Exact (Full Count)O(N) = 10M+ entriesPerfectExpensive sort
Space-Saving ⭐O(K) = 1000 entriesApproximate. Measured error target is within 5%.O(K) reference, O(log K) with heap
Hybrid (Count-Min + Heap + Exact refinement)O(10K)Approximate until exact candidate refinementFast candidate filtering

MapReduce vs Stream Processing

Balancing batch aggregation with streaming analytics provides low latency ranking freshness while retaining an exact historical reconciliation path.

YAML
batch_processing_spark:
  latency: "Hours"
  accuracy: "Exact historical counts after event deduplication"
  architecture: "Full table scans across raw share events with multi-stage aggregation"

stream_processing_flink:
  latency: "Seconds (sub-minute leaderboard freshness)"
  accuracy: "Bounded approximate top-K with exact historical reconciliation"
  architecture: "Continuous stateful stream processing over partitioned Kafka topics"

recommended_hybrid:
  real_time: "Flink computes streaming top-K from pane state every 60 seconds to update Redis leaderboards"
  reconciliation: "Hourly Spark or ClickHouse batch jobs audit exact counts. Large measured drift can trigger an earlier corrective refresh"

Fixed Tumbling Window vs Exponential Decay

Determining how scores decay over time prevents older viral stories from monopolizing leaderboard positions while avoiding abrupt boundary drop-offs.

YAML
fixed_tumbling_window:
  cliff_effect: "Articles drop off abruptly at window boundary (such as T + 1h), causing sudden ranking drops"
  complexity: "Low computational overhead"

exponential_decay:
  formula: "score = shares * exp(-lambda * age)"
  behavior: "Smooth ranking decay over time without abrupt cliff boundaries"
  complexity: "Requires periodic floating-point score recomputations"

recommended_hybrid:
  ingest_tier: "Redis INCR tracks share volume across 1-minute discrete buckets over the last 24 hours"
  scoring_worker: "Flink or scheduled background jobs recompute decay-weighted scores every 5 minutes"
  serving_tier: "Redis sorted sets serve precomputed rankings via ZREVRANGE for sub-50ms reads"

💬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...