Maintain real-time top-K products from events

Read the full interview experience this question came from →

Quick Overview

This question evaluates streaming data aggregation, weighted heavy‑hitter ranking, and the design of efficient data structures and algorithms for real-time top‑K tracking in the Coding & Algorithms domain.

Maintain real-time top-K products from events

Company: Amazon

Role: Software Engineer

Category: Coding & Algorithms

Difficulty: medium

Interview Round: Onsite

Design a data structure/class that ingests a stream of product events and returns the current top‑K products at any time. Events include {timestamp, product_id, event_type} where event_type may be view, add_to_cart, or purchase with configurable weights. Support: update(event), query_top_k(k), and optionally query within a sliding time window W. Specify tie‑breaking, required update and query complexities, and the data structures (e.g., heaps + hash maps, count‑min sketch for heavy hitters). Discuss memory bounds, backfill/replay handling, and concurrency if multiple threads update.

Overview: This question evaluates streaming data aggregation, weighted heavy‑hitter ranking, and the design of efficient data structures and algorithms for real-time top‑K tracking in the Coding & Algorithms domain.

Read the full Amazon Software Engineer interview experience this question came from

Community answers

Answer by Yunan

from collections import defaultdict, deque from dataclasses import dataclass from typing import Dict, List, Tuple, Optional @dataclass class Event: timestamp: int product_id: str event_type: str class TopKProducts: def init(self, weights: Optional[Dict[str, int]] = None, window_size: Optional[int] = None): self.weights = weights or { "view": 1, "add_to_cart": 3, "purchase": 5 } self.window_size = window_size self.score_map = defaultdict(int) # 队列中存储三元组: (timestamp, product_id, earned_score) # 这样过期时可以直接扣除当时计入的分数,避免依赖外部权重变化 self.events = deque() self.max_observed_timestamp = 0 def _apply_event(self, product_id: str, delta: int) -> None: self.score_map[product_id] += delta if self.score_map[product_id] == 0: del self.score_map[product_id] def _evict_expired(self, current_time: int) -> None: if self.window_size is None: return boundary = current_time - self.window_size while self.events and self.events[0][0] <= boundary: _, pid, score_to_deduct = self.events.popleft() self._apply_event(pid, -score_to_deduct) def update(self, event: Event) -> None: # 获取当前事件的权重分 score = self.weights.get(event.event_type, 0) if score == 0: return if self.window_size is not None: # 容错机制:如果事件太旧(已经落后于当前视窗的左边界),则直接丢弃 if event.timestamp < self.max_observed_timestamp - self.window_size: return # 丢弃迟到太久的乱序事件 # 更新已观测到的最大时间戳,防止时钟回拨或乱序破坏边界 self.max_observed_timestamp = max(self.max_observed_timestamp, event.timestamp) # 1. 淘汰过期数据 self._evict_expired(self.max_observed_timestamp) # 2. 将当前事件入队(记录当时确切的分值) self.ev
|Home/Coding & Algorithms/Amazon
Amazon logo
Amazon
Sep 6, 2025
mediumSoftware EngineerOnsiteCoding & Algorithms
5
0

Design a data structure/class that ingests a stream of product events and returns the current top‑K products at any time. Events include {timestamp, product_id, event_type} where event_type may be view, add_to_cart, or purchase with configurable weights. Support: update(event), query_top_k(k), and optionally query within a sliding time window W. Specify tie‑breaking, required update and query complexities, and the data structures (e.g., heaps + hash maps, count‑min sketch for heavy hitters). Discuss memory bounds, backfill/replay handling, and concurrency if multiple threads update.

Submit Your Answer to Earn 20XP

Sign in to leave a comment

Loading comments...