实时热搜榜设计:窗口聚合、抗刷与快照发布
内容平台的实时热搜榜如何设计?
阅读说明与事实边界
“热搜榜怎么做”经常被答成一句话:把话题和分数放进 Redis ZSET,按分数倒序取前 50 个。
ZSET 只能解决最后一步排序,回答不了这些业务问题:
- “热”到底由搜索次数、发帖人数、评论数还是增长速度决定;
- 同一个用户连续搜索 100 次是否算 100 次;
- 一个拥有巨大历史流量的老话题是否永远压住新事件;
- 机器人突然刷高一个词时,榜单怎样止血和回溯;
- 事件晚到、重复、乱序后,窗口分数怎样收敛;
- 算分规则升级时,怎样避免榜单瞬间抖动;
- 运营置顶、风险下榜和自然热度怎样同时表达;
- Redis 丢失后,能否从原始行为和聚合结果重建。
本文以一个综合内容平台为设计案例:用户可以搜索、发布、浏览、评论和转发;平台提供全站热搜榜以及少量频道榜;榜单每隔一个短周期发布快照,展示前 50 个话题和“新、热、沸”等标签。
文中的窗口长度、权重、发布周期和流量都只是设计参数。没有真实埋点分布、压测报告和线上监控时,不写成生产成绩,也不宣称某个公式是所有平台的标准答案。
全文统一使用 scope_id 表示一张榜的隔离范围,而不在各章节分别使用含义模糊的 scope、频道或地域字段:
scope_id = canonical(scope_type, channel_id, region_id, language)
示例:GLOBAL|ALL|ALL|zh-CN
CHANNEL|technology|CN|zh-CN
缺省维度必须写成明确的 ALL,不能用空字符串、null 和缺字段表达三种不同含义。scope_type 说明这是全站、频道还是其他产品榜,后三个维度描述频道、地域和语言;服务端按固定字段顺序编码或计算摘要,禁止客户端自由拼接。一个行为可以按产品配置派生到全站榜和频道榜等多个 scope_id,但每份聚合、候选、快照、缓存和风险动作都只能属于一个确定 scope。window_type 是 scope 内部的特征维度,不等同于对外榜单范围。
一、项目基础:先定义榜单在解决什么问题
1.1 热搜不是行为总量排行榜
行为总量大的话题可能只是常年热门。例如“天气”“某热门游戏”每天都有大量搜索,但用户打开热搜通常想看到“此刻正在发生什么”。
因此热度至少要同时考虑:
- 当前窗口的有效行为量;
- 当前窗口相对历史基线的增长速度;
- 参与用户的去重数量,而不是单纯点击次数;
- 多种行为之间的质量差异;
- 时间衰减,让旧事件自然退出;
- 风险、内容安全和人工运营约束。
1.2 榜单的用户承诺
本文采用以下口径:
| 项目 | 业务口径 |
|---|---|
| 实时性 | 秒级到分钟级更新,不承诺每次点击都立即改变名次 |
| 排名含义 | 平台在指定规则版本下对当前热度的估计 |
| 精确性 | 不承诺对所有原始点击做财务级精确计数,但关键聚合可回放、可解释 |
| 去重 | 同一主体在短时间内的重复行为合并或降权 |
| 下榜 | 自然衰减、风险处置、内容失效或运营规则均可能触发 |
| 置顶 | 与自然排名分开记录并明确审计,不悄悄修改原始热度 |
| 历史 | 保存已发布快照和规则版本,能解释某时刻为何上榜 |
热搜榜是推荐和发现产品,不是财务账本。允许使用近似去重和最终一致,但不能因此丢掉原始证据、规则版本和风险操作记录。
1.3 系统边界
榜单系统不负责决定帖子本身是否违法,也不替代搜索引擎的召回和排序。内容安全服务提供风险结论,话题服务负责把不同文本映射到稳定话题,榜单系统负责聚合、算分和发布。
1.4 非目标
- 不在用户请求时现场扫描所有行为表;
- 不让 App 直接读取内部 Redis ZSET;
- 不把运营置顶混进自然热度分数后失去审计;
- 不用全局唯一的单线程任务处理全站所有话题;
- 不承诺同一毫秒内所有用户看到完全相同的过渡状态;
- 不把“榜单已发布”说成“每个用户已经看到”。
二、业务闭环:一个突发话题怎样进入又退出热搜
以“某地地铁临时停运”为例:用户开始搜索并发布相关内容,平台将别名归并到一个话题,计算有效热度,通过风险检查后发布榜单;事件结束后,新增行为降低,分数随时间衰减并自然下榜。
这条链路中至少有三种不同事实:
- 原始行为事实:谁在什么时间做了什么;
- 窗口聚合事实:某话题在某个窗口有多少有效主体和行为;
- 已发布产品事实:某个榜单版本展示了哪些话题及标签。
不能只保留最终 ZSET,因为它既无法重算,也无法解释运营或风控变更。
2.1 话题生命周期
SUPPRESSED 不等于删除原始行为。风险下榜需要保存原因、操作主体、有效区间和复核结论。解除屏蔽后也不能直接恢复原名次,而应重新按当前窗口和规则计算。
2.2 事件处理状态
原始事件至少一次投递时,重复是正常情况。event_id 去重、窗口聚合的幂等键或上游 Outbox 共同降低重复影响;不能假设消息中间件只交付一次。
三、话题归一:先解决“同一件事有十种写法”
3.1 为什么不能直接用搜索词做 ZSET member
用户可能搜索:
3号线停运
地铁三号线停了
今天地铁3号线怎么了
#地铁3号线临时停运#
直接按字符串计数会把同一事件拆散,也会让恶意用户通过微调文本制造多个榜单条目。
3.2 话题服务输出稳定标识
话题服务把文本标准化并映射为稳定 topic_id,同时保留:
- 标准展示名;
- 别名和标签;
- 频道、地域和语言;
- 关联实体与内容集合;
- 合并、拆分和人工修正历史;
- 当前归一模型或规则版本。
在线热度事件携带 topic_id 与映射版本。模型低置信度时可以先进入“候选词”队列,达到最小门槛后由规则或人工确认,避免垃圾词直接污染榜单。
3.3 合并与拆分
两个话题被确认是同一事件时,不能改个名字后把两个在线计数直接相加。两个窗口可能已消费过相同内容或同一 actor,直接求和会重复计数;映射切换期间还会有新事件和迟到事件并发到达,处理不好又会漏算。
每次合并建立不可复用的 merge_run_id,记录源话题、主话题、目标 mapping_version 和受影响范围。合并意图可以覆盖多个 scope,但交接不能只在父任务上保存一个状态:每个受影响的 scope_id 都建立一条 MERGE_SCOPE_CUTOVER,独立记录期望的 scope context version、新的 mapping_activation_epoch、effective_watermark、交接向量、影子状态引用和交接状态。这样频道榜可以已经切换、地域榜仍在修复,恢复程序不会把“部分完成”误判成全局完成。
交接位点只取自契约校验之后、话题归一之前的稳定事件日志。该日志中的 source_partition + source_offset 一旦产生就不再因 topic 或 scope 改变;后续归一、scope 展开和重分区必须一直携带这两个字段。真正决定新旧链路各自负责哪些事件的,不是事件时间,也不是归一后的 topic 分区 offset,而是这份归一前日志的交接位点向量:
cutover_next_offset[p] = 归一前稳定日志分区 p 中第一条必须按新映射处理的 offset
也就是:offset < cutover_next_offset[p] 归旧映射基线,offset >= cutover_next_offset[p] 归新映射链路。effective_watermark 仍用于判断窗口是否允许修改,但不能兼任消费所有权边界。
影子状态不能只重算几个数字。它要同时重建受影响窗口的事件去重集合、actor 去重或近似基数状态、局部聚合状态和源位点,否则切换后同一 actor 在源话题和主话题各出现一次时仍可能重复计数。为了缩短 barrier 处的暂停,可以先从一致检查点构建影子状态并持续追增量,等它接近在线水位后再插入 barrier;暂停期间只补最后一小段并做摘要校验。
每个 scope 的交接 CAS 是该 scope 唯一的线性化点。它必须在同一控制事务中锁定对应 RANKING_SCOPE_CONTEXT,校验期望 context version 和完整交接向量,推进 mapping_activation_epoch、登记新的活动 mapping_version、激活影子状态并把该 scope cutover 标为 ACTIVE。CAS 之前失败,该 scope 仍认旧映射;CAS 之后恢复,只能从持久 scope 上下文认新映射。父 merge run 只有在所有 scope cutover 都进入终态后才能结束,不能用父状态覆盖子 scope 的真实进度。
归一工作结果必须携带 source_partition + source_offset + scope_id + mapping_activation_epoch。旧 normalizer 可能在 barrier 后才吐出旧 epoch 结果,目标聚合器不能直接接收:它根据 scope 上下文拒绝旧 epoch,并按稳定源位点把原事件重新送入当前映射链路;source_partition + source_offset + scope_id 的幂等记录阻止重路由与原结果双算。这样归一后的 topic 即使发生重分区,也不会失去可重放的事件所有权边界。
事件时间只在交接所有权确定之后发挥作用。例如一条源 offset >= cutover_next_offset[p]、但 event_time <= W 的迟到事件,仍然只归新映射链路:窗口处于 CLOSED 且在允许迟到范围内就修正主话题,窗口已经 FINAL 就进入隔离审计或授权回放,不能送回旧合并任务。这样同一条事件不会因为“晚到”同时被在线链路和补偿链路处理。
影子构建和交接记录的幂等键至少包含 merge_run_id + scope_id + window_type + window_start + metric_type + topic_id + compensation_version。每次写入既校验对应聚合行版本,也校验该 MERGE_SCOPE_CUTOVER 的完整 cutover_next_offset 向量和 mapping_activation_epoch。历史已发布快照仍保存当时的话题身份,另记录事后归并关系;只有新快照使用新映射,且同一主话题只能出现一次。
如果一个大话题后来需要拆成两个事件,历史行为无法凭空精确拆分。应基于内容和查询重新归类可识别事件,其余保留在原话题并标注估计边界。
四、数据模型:保存行为、聚合和发布三层证据
这里故意不让 RAW_EVENT_PARTITION 直接关联某个窗口版本。原始分区只描述稳定日志中的事件所有权;某次计算实际读到哪里,由不可变的 SOURCE_OFFSET_CUT 及其 SOURCE_OFFSET_CUT_ITEM 固化,窗口版本再引用该 cut。于是同一原始分区可以出现在多个时间点的 cut 中,一个窗口版本也可以引用覆盖多个原始分区的完整向量,不会被错误建模成“一条分区记录直接产出一个窗口版本”。
4.1 为什么需要榜单快照
如果客户端每次请求都读取一个正在被更新的 ZSET,第一页和下一次刷新可能来自不同计算时刻,运营也无法回答“18:30 的第 3 名是什么”。
榜单服务生成不可变 snapshot_id,写全量榜项后再原子切换当前版本指针。历史快照按保留策略归档,用于产品复盘、申诉和规则评估。
4.2 关键唯一约束和索引
- 原始事件以
source + event_id唯一,或由可重放日志位点证明处理范围; TOPIC_WINDOW_METRIC是窗口 head,以topic_id + scope_id + window_type + window_start + metric_type唯一,只保存当前代次和版本指针;每次聚合先追加不可变的TOPIC_WINDOW_METRIC_VERSION(aggregate_key, aggregation_generation, aggregate_version),其中固化指标值、独立输入水位、源 offset 向量引用和值摘要,再以expected active_generation + expected current_aggregate_version双重 CAS 推进 head;- 每个候选输入清单以不可变
input_manifest_id标识,清单项必须用aggregate_key + aggregation_generation + aggregate_version外键引用确切版本;三者与source_offset_vector_ref + value_digest一并参加清单摘要,不能只抄 head 上一个随时会前进的版本号,也不能把多行输入压成一个可比较的“最大聚合版本”; RUN_PARTITION_CUT在 run 创建时为每个预期候选分区插入一行PENDING,candidate_partition_id + expected_source_cut_ref此后不可改。分区只能用一次条件写把topk_object_ref + topk_item_count + topk_digest补齐并推进为COMPLETED;进入BUILDING后所有 cut 行只读。预期集合摘要覆盖分区和各自 cut,结果集合摘要还覆盖每个分区的结果对象、条数、摘要和完成状态,不能只对“已经上报的行”求摘要;- 每轮 run 在开始计算前先冻结
RUN_INPUT_MANIFEST集合。input_manifest_set_digest按topic_id + input_manifest_id + manifest_digest稳定排序计算;每个条目最终只能进入CANDIDATE或带原因的REJECTED,从而能证明没有漏处理某个输入; - 候选以
scope_id + calculation_run_id + topic_id隔离,写候选与确认父 run 仍为BUILDING放在同一事务;数据库拒绝对SEALEDrun 新增、更新或删除候选。发布器只读取一个已封口的正式 run;影子规则和回放任务另放独立 namespace,并携带输入清单与三类 activation epoch; - 快照项以
snapshot_id + rank_no唯一; - 同一快照内
topic_id唯一,防止合并话题重复上榜; - 风险操作按
scope_id + topic_id + effective_from索引;跨所有榜的动作使用明确的全局 scope,再由查询规则与局部 scope 共同求值; RANKING_SCOPE_CONTEXT以scope_id唯一,保存当前激活的规则、窗口配置、映射、scope 风控纪元及单调控制版本;规则回滚也只能继续增加 activation epoch;RANKING_PUBLISH_POINTER以scope_id唯一,是已发布版本的数据库权威记录;Redis 只缓存它;- 当前榜单指针更新时条件锁定并校验 scope 上下文版本,同时 CAS 指针版本;发布事务本身不推进 context version。规则、映射和风险控制事务也锁同一条 context 行,因此它们与发布之间有明确先后,旧生成任务不能覆盖新快照。
五、技术架构与模块边界
5.1 各组件只承担一类责任
| 组件 | 负责 | 不负责 |
|---|---|---|
| 业务服务 | 记录搜索、发布和互动事实 | 现场计算全局热搜 |
| 话题归一 | 文本到稳定 topic 的映射 | 决定最终排名 |
| 流计算 | 时间窗口、去重和聚合 | 永久保存所有原始载荷 |
| 算分服务 | 规则版本下的自然热度 | 内容合规最终裁决 |
| 风险校正 | 屏蔽、降权和异常主体影响 | 偷偷改写原始指标 |
| Scope 控制上下文 | 分配规则、映射、风控 activation epoch,保存当前允许发布的上下文 | 保存全部窗口明细或直接服务榜项排序 |
| 快照发布器 | 锁定但不修改 scope context,生成稳定快照并 CAS 数据库权威指针 | 在客户端请求内重新算分,或推进规则、映射和风险 epoch |
| Redis | 服务当前和短期历史快照 | 作为唯一行为事实源 |
| 边缘风险过滤 | 在 CDN 命中后仍执行紧急 denylist,或强制绕过不安全缓存 | 替代正式风险计算和审计 |
5.2 同步与异步边界
用户搜索或互动的主请求只需可靠记录自己的业务事实或 Outbox,不同步等待热搜计算。榜单异步更新,允许短暂延迟。规则激活、映射交接和风险动作属于低频控制事务,必须先持久推进 scope context,异步计算结果才能据此获得发布权。风险紧急下榜可以走独立高优先级控制通道,先增加带 risk epoch 的边缘与查询层拦截,再生成新快照;操作必须落审计库并最终反馈到正常计算链路。只清 CDN 而没有命中后过滤,无法形成安全承诺。
六、核心技术实现:事件接入怎样处理重复、乱序和迟到
6.1 统一事件契约
{
"eventId": "01K...",
"eventType": "SEARCH_EXECUTED",
"eventTime": "2026-08-19T20:00:01+08:00",
"receivedAt": "2026-08-19T20:00:02+08:00",
"actorKey": "privacy-safe-key",
"contentOrQueryId": "q_1001",
"topicHints": ["t_9001"],
"scopeDimensions": {
"channelId": "local-news",
"regionId": "CN-SH",
"language": "zh-CN"
},
"sourceService": "search-service",
"schemaVersion": 3,
"mappingVersion": 42,
"riskContextVersion": 12
}
actorKey 应是满足隐私和风控需要的稳定主体键,不在热搜日志中复制手机号等敏感信息。事件时间用于窗口归属,接收时间用于判断延迟和排查链路。事件里的频道、地域和语言是输入维度,不允许上游直接指定最终榜单 Key;mappingVersion 表示 topic hint 所依据的映射版本,归一服务仍需按当前有效映射核验。
契约校验通过后,事件先写入归一前稳定日志。source_partition + source_offset 是日志系统分配的信封字段,不由上游伪造;从归一输出、scope 展开、重分区、窗口版本到候选输入清单都必须继续携带或引用它。话题合并的 barrier 只使用这组稳定坐标,不能使用 topic 已改变之后的下游分区 offset。
归一阶段输出稳定 topic_id + mapping_version + mapping_activation_epoch,并按榜单配置把一个事件展开成一个或多个确定的 scope_id,例如同时进入 GLOBAL|ALL|ALL|zh-CN 和 CHANNEL|local-news|CN-SH|zh-CN。每份展开记录都保留原 event_id、归一前 source_partition + source_offset 和 scope 派生规则版本,使重复展开可去重、旧 epoch 输出可拒绝并按源位点重路由、scope 归属可重放。其下游聚合键固定为:
topic_id + scope_id + window_type + window_start + metric_type
6.2 正常接入时序
榜单事件投递失败不能拖垮搜索主链路,但 Outbox 或可靠埋点缓冲必须让失败可见、可重试。客户端直传埋点适合曝光等统计,关键的发帖和搜索事实优先由服务端生成。
6.3 水位和允许迟到
不能用“当前机器时间减一分钟”代表所有输入都已到齐。对每个参与该 scope_id + window_type 的数据源,先按分区计算:
partition_watermark = max_observed_event_time - configured_out_of_orderness
source_watermark = min(non_idle_partition_watermark)
window_watermark = min(required_source_watermark)
configured_out_of_orderness 按数据源真实乱序分布配置。一个分区超过 idle_partition_timeout 没有事件和心跳,才可暂时从最小值计算中移除,同时把该 source/scope 标成数据不完整;不能因为低流量分区短时没消息就随意判 idle。必需数据源不完整时,发布器按策略冻结新榜或切到明确的降级规则版本。
空闲分区恢复后先进入 CATCHING_UP:从已记录位点读取积压,早于当前 watermark 的记录按允许迟到规则处理,不能把全局 watermark 倒退。只有该分区追到当前安全水位并通过位点连续性检查,才重新纳入非 idle 分区集合。这个过程需要记录 idle 起止时间、恢复位点和被送入迟到侧输出的数量。
窗口按以下状态收敛:
OPEN接收正常事件;CLOSED已产生可供算分的结果,但仍接收允许迟到范围内的修正,修正只能影响后续快照;FINAL的正式聚合不可再写,超时事件进入审计或隔离回放;PURGED前必须保存输入水位、源位点摘要、最终聚合版本、检查点或归档位置以及对账结果。
去重状态的保留期不得短于“最大允许迟到时间、最大自动重放范围、消息最大重投窗口”三者中的最大值。清理顺序必须是:窗口进入 FINAL,保存检查点与归档证据,确认源位点和聚合对账,再清理去重状态和窗口状态。否则历史消息在去重集合过期后自动重投,会被当成新事件重复累计。
没有必要为了追求理论精确而不断改写用户已经看到的历史榜单。已发布快照保存当时可用事实,离线重算结果作为评估版本另存;若要影响后续正式榜,必须走隔离验证和授权发布。
6.4 分区与顺序
归一和 scope 展开后,事件可以按 topic_id + scope_id 分区,使同一话题在同一榜单范围内的窗口聚合逻辑有序。话题归一发生前,可按查询或内容哈希进入预处理分区,再按 topic 和 scope 重分区。热门话题会形成分区热点,需要在聚合阶段做本地预聚合或话题加盐分片,但盐值只是局部计算维度,最终仍合并到完整业务聚合键,不能遗漏 scope、窗口或指标类型。
七、窗口聚合:既看当前量,也看增长速度
7.1 多尺度窗口
仅用一分钟窗口会剧烈抖动,仅用一天窗口又无法发现突发事件。可以同时维护:
- 短窗口:发现突发增长;
- 中窗口:判断热度是否持续;
- 长基线:衡量正常历史水平;
- 衰减累计:让热度平滑退出。
窗口长度由产品节奏和真实流量决定。体育比赛和财经资讯可能需要更短窗口,知识社区可以更长。
7.2 有效指标而不是原始点击
每个话题按窗口聚合:
| 指标 | 价值 | 常见风险 |
|---|---|---|
| 去重搜索主体数 | 主动意图强 | 刷搜索词、自动补全误触 |
| 去重发帖主体数 | 讨论供给 | 搬运、批量机器账号 |
| 有效内容量 | 事件信息丰富度 | 垃圾内容堆量 |
| 评论和转发主体数 | 传播与参与 | 互刷和营销群控 |
| 点击后停留或深度阅读 | 内容质量 | 采集成本和隐私边界 |
| 举报、隐藏和负反馈 | 风险与低质信号 | 恶意举报 |
“去重主体数”可以用精确集合或 HyperLogLog 等近似结构。前 50 候选话题需要更精细审计时,可以对候选集二次精算;长尾话题用近似聚合降低成本。
7.3 分层聚合避免热点打满单节点
局部聚合只减少网络和写放大,不改变去重语义。局部结果携带 shard、输入水位、源位点摘要和局部版本,合并器按完整业务键收齐或按明确缺失策略推进。若 actor 可能落在多个局部分区,不能简单相加精确去重数;可使用可合并的近似基数结构,或按 actor 稳定分区后再聚合。
八、算分:公式必须可解释、可版本化
8.1 一个示意框架
不把某个固定公式写成标准答案,可以把自然热度拆成:
自然热度
= 当前有效行为强度
+ 相对历史基线的增速
+ 多样性与内容质量奖励
- 时间衰减
- 重复主体和异常流量惩罚
不同指标先做尺度归一,防止原始数量大的搜索完全吞掉评论和发帖信号。权重、归一方式、阈值、衰减参数都属于 score_rule_version。
8.2 防止老话题霸榜
假设话题 A 当前搜索量很大但与历史相同,话题 B 总量较小却在几分钟内增长十倍。热搜应给 B 足够机会。可以结合:
- 当前值与同星期、同时段历史基线的比值;
- 短窗口相对中窗口的斜率;
- 新话题冷启动奖励,但设置最低有效主体门槛;
- 在榜时长衰减或重复曝光疲劳;
- 退出阈值低于进入阈值,减少边界反复抖动。
8.3 进入与退出的滞回
例如候选进入需要同时满足最低有效主体和进入分数;已经在榜的话题只有连续若干快照低于退出阈值才下榜。这是产品稳定性设计,不是篡改热度。
8.4 同分排序
同分时使用稳定规则,例如增长速度、首次达到门槛时间和 topic_id。不能每次由无序 Map 遍历结果决定,否则榜单会在完全相同数据下抖动。
九、防刷和内容安全:不要只按 IP 限流
9.1 三层防护
第一层处理缺字段、异常时间、伪造来源和明显重放;第二层处理同账户、设备、IP 或匿名主体的短期重复;第三层识别大量新账号同步行为、相似内容矩阵、异常地域和转发图谱。
9.2 原始值与校正值分开
窗口和候选中同时保存:
- 原始事件数;
- 基础去重后的有效值;
- 风控识别的可疑数;
- 风险调整后的参与算分值;
scope_id、话题风险操作版本、scope 风控纪元和风险模型版本。
风险动作必须声明作用域:局部动作写具体 scope_id,需要覆盖所有榜时写明确的 GLOBAL|ALL|ALL|ALL 并由查询层同时求值,不能靠缺少 scope 字段暗示全局。解除屏蔽不是删除旧记录,而是产生更高的 risk_action_version。这样才能解释“明明搜索很多为什么没有上榜”,也能在风控误判后重新计算。不要把可疑事件物理删除到无证据可查。
单个话题的 risk_action_version 不能直接充当整张榜的风险版本:话题 A 的第 8 次操作与话题 B 的第 3 次操作没有可比较关系。每个 scope 还要维护单调递增的 risk_epoch。新增、解除或到期一个会改变榜单结果的局部风险动作时,在同一控制事务中写动作并推进该 scope 的 risk_epoch;切换风险模型也推进它。全局动作另有 global_risk_epoch,可由 GLOBAL|ALL|ALL|ALL 这条权威 scope 控制记录分配。边缘和查询层立即求全局与局部动作并集,各 scope 的控制投影随后记录已吸收的全局 epoch;发布事务还要校验这个值等于全局权威行的当前 epoch。在一个 scope 尚未吸收最新全局 epoch 前,发布器不得生成声称“风险上下文最新”的新快照。
effective_until 不是“到了墙上时间就自动解除”的读取规则,而是到期任务最早可以提交 EXPIRED 控制动作的时刻。动作只有在到期任务锁定 scope context、写入更高动作版本并推进 scope 或 global risk epoch 后才失效;任务延迟时宁可继续保守屏蔽,也不能静默放行而不换 epoch。任务使用动作 ID 幂等执行并可重试,快照和 CDN 的有效期不得跨过当前最早到期时间。这样重复执行、调度延迟和全局动作到期都有持久证据。
候选分不是只有一个可被覆盖的浮点数。一次算分通常读取多个窗口和多个指标,因此先生成不可变输入清单:
input_manifest_id -> [
(window_1m, search_uv, generation=g7, aggregate_version=18, value_digest=...),
(window_5m, post_uv, generation=g7, aggregate_version=11, value_digest=...),
(window_1h, baseline, generation=g6, aggregate_version=42, value_digest=...)
]
清单同时保存输入水位、源 offset 向量引用和窗口配置版本。每一项用 aggregate_key + aggregation_generation + aggregate_version 引用不可变的 TOPIC_WINDOW_METRIC_VERSION,该版本行内保存实际指标值和值摘要;窗口 head 即使随后从 g7/v18 前进到 g7/v19,已冻结 run 仍能按原三元组重现结果。恢复代次可能从 g7/v19 切到 g8/v3,所以版本号只能在同一 generation 内比较;一个脱离 generation 的标量 aggregate_version=42 既不能确定引用哪一行,也不能表达多窗口输入的偏序关系。
候选的完整身份可以摘要为:
score_version = digest(
input_manifest_id,
scope_context_version,
rule_version, rule_activation_epoch,
mapping_version, mapping_activation_epoch,
risk_epoch, risk_model_version, global_risk_epoch,
calculation_run_id, candidate_revision
)
score_version 只是身份和校验摘要,不能按字符串或哈希大小判断新旧。候选写入必须逐项验证 scope 上下文:scope_context_version、活动窗口配置和三个 activation epoch 都与当前控制记录一致;输入清单引用的每个不可变版本存在、值摘要相符,并且都属于本 run 固定的 source-offset cut;同一 calculation_run_id 内只接受更高 candidate_revision。**不要求窗口 head 停止前进,也不因 head 出现 v19 就让已经封口的 v18 run 失去可重复性。**新水位生成新的清单和新 run,不把两个摘要强行比较。
风险请求也要携带 scope_id + scope_context_version + input_manifest_id + rule_activation_epoch + mapping_activation_epoch + risk_epoch + global_risk_epoch + calculation_run_id。延迟响应返回时只要任一项已经变化,CAS 就失败并重新计算,不能覆盖较新结果。影子规则和历史回放使用隔离 namespace;即使它的数字或哈希看起来更大,也没有正式发布权。
9.3 紧急下榜
内容安全确认需要立即下榜时:
- 在控制事务中写入带
scope_id和递增risk_action_version的风险操作,同时推进该 scope 的risk_epoch; - 把带 epoch 的高优先级 denylist 推送到 CDN 命中之后仍会执行的边缘过滤层,并同步到查询层;边缘尚未确认最新 epoch 时绕过榜单缓存回源,不等待下一张榜;
- 候选按新的风险 epoch 重算,发布器生成不包含该话题的新快照;
- 清理所有受影响 scope 的 CDN 对象;全局动作展开受影响 scope 清单,purge 失败时由边缘 denylist 继续过滤,缓存最大安全 TTL 作为最后边界;
- 记录操作人、规则、原因、动作版本、风险 epoch 和生效时间;
- 正常计算继续保留自然分,用于复核,不自动对外展示。
紧急通道必须可审计且有显式到期事务,避免一次临时屏蔽永久遗留。解除屏蔽或到期同样生成新动作版本、推进 epoch、触发候选重算和新快照,不能直接删掉 Redis 屏蔽 member;全局动作和局部动作的并集仍在当前 scope 下求值,避免频道 A 的操作误伤频道 B。若平台没有 CDN 命中后的可编程过滤能力,安全口径只能选择短 TTL 并在紧急状态下绕过 CDN,不能声称“查询层一定即时拦截”。
9.4 运营干预
置顶、保底曝光、活动话题位和自然热搜最好分成不同产品位。如果必须混排,快照项同时保存自然分、调整类型和最终位置。运营不能直接把数据库的自然分改成一个更大数字。
十、候选榜与快照发布
10.1 为什么分两阶段
流处理不断更新候选分数,而对外榜单需要稳定版本。候选不能只按 scope_id + topic_id 原地覆盖,否则本轮只更新了一半时,发布器可能把上一轮和本轮拼成一张榜。每轮正式计算使用独立 calculation_run_id:
- 创建 run 时先冻结候选生产阶段的
source_offset_vector_ref和预期分区集合,按candidate_partition_id写RUN_PARTITION_CUT(PENDING),保存每个分区可验证的expected_source_cut_ref。expected_partition_set_digest按(candidate_partition_id, expected_source_cut_ref)稳定排序计算,不能只摘要分区 ID;run 此时是COLLECTING; - 每个分区只计算自己固定 cut 以内的局部 Top K,把不可变的
topk_object_ref + topk_item_count + topk_digest一次性提交后,才条件推进为COMPLETED。重复提交只有内容完全相同才按幂等成功,整个分区未上报不能靠其他分区的完整结果掩盖; - 协调器确认
COMPLETED行数等于预期数、分区 ID 一个不少且每个 cut 与 run 的源 offset 向量一致,再按(candidate_partition_id, expected_source_cut_ref, topk_object_ref, topk_item_count, topk_digest, completion_status)稳定排序计算partition_result_set_digest。校验通过后才合并所有局部 Top K、写入完整RUN_INPUT_MANIFEST,保存input_manifest_count + input_manifest_set_digest,再 CASCOLLECTING -> BUILDING; BUILDING阶段只允许处理这个冻结集合。每个输入条目必须落成CANDIDATE或带原因的REJECTED,不能靠“最后大约有 N 条”判断完成;- worker 写候选、写 result digest 和把条目标终态时,在同一事务内锁定父 run 并确认仍为
BUILDING; - 封口事务锁定 run,验证期望输入集合一个不少、一个不多,所有条目均为终态,候选 topic 集合与
CANDIDATE条目精确相等,再写candidate_count + candidate_set_digest并 CAS 为SEALED; BUILDING后分区 cut 和输入集合不可修改;SEALED后数据库再拒绝候选的新增、修改或删除,迟到修正只能生成新 run;- 发布器只从一个
SEALED正式 run 读取 Top K,影子、回放和未封口 run 没有发布资格。
如果某个候选分区长期不可用,默认做法是让原 run 停在 COLLECTING 并继续服务上一张快照,原 run 不能就地删掉缺失分区。业务确实允许残缺数据发布时,应依据已审批的 degraded_policy_version 创建一个新 run,把“纳入分区集合、排除分区集合、适用 scope 和有效期”一起固化到降级策略及摘要中,再为新预期集合创建 cut;新 run 也必须把自己的全部预期分区收齐后才能生成 manifest,并写 data_completeness_status=DEGRADED。这样审计才能区分“全量 Top K”与“明确少了哪些分区但仍获准发布”。
这样发布器按短周期读取的是一套冻结候选,再做二次校验和去重,生成前 N 快照。
快照状态至少区分 PREPARED -> READY -> PUBLISHED,竞争失败、版本过期或校验失败进入 REJECTED。READY 只说明内容完整,不代表产品已经发布;数据库必须在同一事务中完成 RANKING_PUBLISH_POINTER 的 CAS 和快照状态改为 PUBLISHED,任一步失败都回滚。写快照内容和切换权威指针必须分开,客户端不会读到半张榜,也不会在 Redis 丢失时误选一个从未发布成功的 READY 快照。
每个 scope_id 只有一条数据库权威发布指针和一条 scope 控制上下文。前者回答“用户当前读哪张已发布榜”,后者回答“现在允许哪个规则、映射和风控上下文发布”。两者不能混成 Redis 里的一条可随意覆盖记录。基线把 context、pointer 和 snapshot 状态放在同一事务数据库中:发布事务先条件锁定 context 行并验证版本,不增加 context version,再 CAS 指针并发布快照;规则、映射和风险事务也锁同一 context 行并负责推进它的版本。若它们被拆到不同数据库,就必须另设单一发布仲裁服务,不能把两次本地 CAS 冒充原子提交。发布事务至少同时验证:
- 指针的
scope_id与快照 scope 完全一致; - 读取时的
version和publication_epoch仍是当前值; - 读取时的
RANKING_SCOPE_CONTEXT.version没有变化; - 快照的
rule_version + rule_activation_epoch等于当前活动规则上下文; - 快照的
window_config_version等于当前活动窗口配置; - 快照的
mapping_version + mapping_activation_epoch等于当前活动映射上下文; - 快照的
risk_epoch + risk_model_version + global_risk_epoch等于当前活动风控上下文,并且其中的 global epoch 等于全局权威 scope 行; calculation_run_id仍是SEALED,run 的scope_context_version、输入水位、分区结果摘要、input_manifest_set_digest、输入条目数、候选条目数和candidate_set_digest都自洽;快照保存的两类 digest 和两类 count 与该 run 完全一致;- 每个快照项引用的输入清单都在该 run 的冻结集合内,清单项引用的不可变聚合版本、值摘要和源 offset cut 一致,
candidate_set_digest能覆盖所有榜项; input_watermark不低于当前已发布水位,且窗口配置和必需数据源完整性满足本次规则;calculation_run_id属于正式在线 run,或携带经过审批、限定 scope、水位和三类 activation epoch 的回放 promotion 授权。
publication_epoch 由数据库在 CAS 成功时单调分配。任一风险动作、规则激活或映射切换发生在“读候选”与“发布”之间,都会推进 scope context 或对应 epoch,使本次 CAS 失败。Redis 中的 current pointer 只是数据库发布指针的缓存,写入时也按 publication epoch 拒绝旧值;即使旧发布器最后才恢复,也无法用旧规则、旧水位或旧 run 抢回当前榜。
10.2 发布前检查
- 候选话题仍处于可发布状态;
- 同一主话题没有通过别名重复出现;
- run 的预期候选分区、各分区 cut 与 Top K 摘要、输入 manifest 集合、终态 disposition、候选集合和各层集合摘要精确闭合,封口后没有写入;
- scope、输入水位、窗口配置、三类 activation epoch 和各榜项输入清单固定且彼此匹配;
- 同一榜项所需的多窗口、多指标没有缺行,也没有拿一个最大
aggregate_version冒充输入向量; - 分数不是 NaN、负无穷或异常跳变;
- 榜项数量、频道配额和运营位符合约束;
- 与上一快照相比变化过大时触发保护或人工检查;
- 每个榜项生成简短解释摘要,供内部审计。
10.3 查询接口
GET /rankings/hot-search?scopeType=channel&channelId=technology®ionId=CN&language=zh-CN&snapshot=latest
查询服务按与计算侧相同的规范化规则生成 scope_id;字段组合不合法时直接拒绝,不能回退到全站榜。返回 scope_id、snapshot_id、publication_epoch、generated_at、input_watermark、rule_version、rule_activation_epoch、risk_epoch、items。客户端使用 ETag 或 snapshot_id 做条件请求。
缓存键必须保留 scope 隔离,例如:
ranking:pointer:{scope_id}
ranking:snapshot:{scope_id}:{snapshot_id}
ranking:emergency-block:{scope_id}
CDN cache key = canonical path + scopeType + channelId + regionId + language + snapshot + cache_policy_version
CDN 可以短缓存榜单,但紧急屏蔽必须按受影响 scope_id 主动失效,并在 CDN 命中之后再经过边缘 denylist 过滤;请求携带的边缘 risk epoch 落后时必须绕过缓存回源。全局动作要展开并失效所有受影响 scope,purge 失败由过滤层和最大安全 TTL 承担边界。不能只用 /rankings/hot-search?snapshot=latest 做缓存键,否则全站、频道、地域或语言榜可能串读;也不能只有源站查询过滤却让 CDN 直接把旧完整响应返回给用户。
10.4 个性化与全站榜不要混淆
全站热搜是共同快照;个性化趋势可以在全站候选上按用户兴趣二次排序,但产品必须明确它不是所有人相同的“全站排名”。两者的缓存键、解释和合规要求不同。
十一、规则升级如何避免榜单瞬间失控
11.1 规则不可原地覆盖
每次权重、阈值、衰减或风险因子变化生成新 rule_version。变更包含:
- 配置内容和摘要;
- 创建人、审核人和变更原因;
- 适用榜单范围;
- 计划生效时间;
- 回滚版本;
- 离线回放和影子结果。
rule_version 表示不可变配置内容,rule_activation_epoch 表示某个 scope 第几次激活规则,两者不能合并。每次上线、灰度扩大、回滚或再次启用同一份旧配置,都在 scope 控制事务中分配一个更大的 activation epoch。例如:
epoch 17 -> rule-v2
epoch 18 -> rule-v3
epoch 19 -> rule-v2 // 回滚到旧内容,但不是回到旧时代
如果回滚时把 epoch 从 18 改回 17,早已暂停的 rule-v2 / epoch 17 任务可能在恢复后再次满足发布条件,这就是 ABA。activation epoch 只能前进;旧任务携带的 epoch 即使 rule_version 相同也必须被拒绝。
11.2 影子计算和差异评估
新旧规则必须固定同一 scope_id、输入水位、窗口配置、候选输入清单、映射 activation epoch 和风险 epoch 比较,不能一个用实时数据、另一个用昨天数据。正式切换在数据库中 CAS RANKING_SCOPE_CONTEXT:写入目标 active_rule_version,同时把 rule_activation_epoch 和 context version 各加一。回滚使用相同步骤,只是目标规则内容换成旧 rule_version,绝不把 epoch 调小或复用旧计算 run;影子数据继续保留在隔离 namespace。
11.3 模型版本
如果使用机器学习做质量或风险校正,模型特征、训练数据时间、推理版本和阈值同样要记录。模型分不能替代可解释的基础指标,至少内部需要知道某话题被压低来自哪个风险或质量信号。
十二、可靠性保证与不保证
| 场景 | 系统保证 | 不保证 | 恢复证据 |
|---|---|---|---|
| 行为事件被业务服务接受 | 关键事件进入业务事实或可靠 Outbox | 不保证立即体现在热搜 | Outbox 状态、事件位点和聚合水位 |
| 事件重复投递 | 聚合按事件或业务键收敛,重复可检测 | 不保证中间件只交付一次 | 去重记录、窗口版本和重放报告 |
| 事件迟到 | 允许范围内修正后续快照 | 不保证无限回改历史榜单 | watermark、迟到队列和离线评估 |
| Redis 榜单缓存丢失 | 按 scope 从数据库权威指针恢复已发布快照 | 不保证恢复期间保持原延迟 | 发布指针、publication_epoch 和快照存储 |
| 风控服务超时 | 按预先定义的风险策略降级,旧响应不能越过 scope 风控纪元 | 不保证所有话题都继续发布 | risk epoch、动作版本、超时指标和人工入口 |
| 发布器宕机 | 数据库权威指针指向的上一完整快照继续服务 | 不保证 READY 快照自动发布 | PUBLISHED 快照和权威指针 |
| 查询成功 | 返回一个明确版本的榜单 | 不代表每个用户都已看到 | 请求日志、snapshot_id 和 CDN 状态 |
12.1 Redis 故障
候选榜 Redis 丢失时,先固定 RANKING_SCOPE_CONTEXT.version 和三类 activation epoch,再从同一 source-offset cut 下的不可变窗口版本生成新的输入清单和 RUN_INPUT_MANIFEST,在隔离候选 namespace 重算并逐榜项对账;上下文在重建期间变化就废弃本批结果重来。对外快照 Redis 丢失时,查询数据库 RANKING_PUBLISH_POINTER,只回填该指针指向的 PUBLISHED 快照及其 publication_epoch。任何“最近 READY”都可能是 CAS 失败或尚未获准发布的结果,绝不能作为恢复依据。若数据库暂时不可用,就继续服务进程内仍可验证的上一已发布快照或明确降级,不现场扫描候选数据生成新榜。
12.2 流处理状态损坏
窗口状态应有检查点和源事件位点。恢复不能把 aggregate_version 塞进 head 的定位键:在线 head 始终用 topic_id + scope_id + window_type + window_start + metric_type 定位。不可变版本则必须以 aggregate_key + aggregation_generation + aggregate_version 唯一;不同 generation 都可以从自己的 v1 开始,互不碰撞,输入清单也用同一个三元组引用。
为避免恢复任务与在线流处理互相覆盖,恢复先在新的 aggregation_generation 中从一致检查点重放,追加不可变版本行并核对各自的 source offset 向量。在线 worker 每次推进 head 都必须同时匹配 expected active_generation + expected current_aggregate_version;恢复追到固定 source cut、完成五字段全键对账后,也用同样的双重 fence 一次性切换 head。切换成功后,旧 generation worker 即使稍后恢复并携带一个数值更大的 version,也会因 generation 不等而失败,只能停止或重新绑定当前代次;切换失败则废弃恢复代次。不能缩写成 topic_id + window,否则频道、窗口长度和指标类型会互相覆盖;也不能既恢复旧状态又从最新位点消费,导致中间事件永久跳过。
例如 head 已从 g7/v101 切到恢复代次 g8/v3,暂停的旧 worker 此时才带着 g7/v102 恢复。若 CAS 只比较数字版本,102 > 3 会错误地把 head 抢回旧状态;双重 fence 会先发现 expected active_generation=g7 与当前 g8 不同而拒绝。g7/v102 可以作为旧代不可变证据保留,却不能成为活动 head,也不能被新 manifest 冒充 g8/v102 引用。
12.3 部分数据源中断
若搜索事件正常但评论事件中断,继续按原公式发布可能系统性偏向搜索。监控每个源的事件速率和延迟;超过阈值时冻结新榜、切换到明确的降级规则版本,或在榜单上使用上一稳定快照。选择取决于数据源的重要性和产品口径。
12.4 重放不能污染正式榜单
历史回放写入隔离命名空间和新的 calculation_run_id。验证完成并不自动获得发布权:若结果需要影响后续正式榜,必须创建 promotion 授权,明确审批人、目标 scope_id、输入水位范围、输入清单摘要、规则/映射/风险 activation epoch、有效期、期望 scope context version 和期望指针版本。补数任务不能直接覆盖当前候选榜,也不能修改历史已发布快照。
十三、分区、热点与容量设计
13.1 先给出容量输入
容量设计至少需要:
各事件类型峰值 EPS
平均和 P99 事件字节数
有效 topic 数及长尾分布
单话题最坏热点比例
窗口数量、长度和允许迟到范围
去重结构的每窗口内存
候选 TopK 大小与发布频率
原始事件保留时间和压缩率
故障恢复时要求的追赶时间
只有“日活多少”无法推导流处理吞吐;只有“QPS 多少”也无法推导窗口状态内存。
13.2 热点话题
突发事件可能占全站大部分事件。如果按 topic 单分区,热点会拖慢同分区其他话题。可以:
- 接入端在每个
scope_id内按topic_id + shard_suffix做局部聚合; - actor 稳定路由以保留去重语义;
- 用可合并草图统计去重主体;
- 热点 topic 使用独立计算资源或动态拆分;
- 合并层实施背压,不能让候选榜写入无限堆积。
13.3 候选集控制
不需要对数百万长尾话题每个发布周期全排序。流计算先应用最低主体和质量门槛,维护分区 Top K,再合并出全局 Top K。正式发布器只精算较小候选集,并检查风险和稳定性。
13.4 背压与降级顺序
- 暂停非关键解释特征和低价值曝光事件;
- 扩大局部聚合批次,降低候选榜更新频率;
- 保留搜索、发帖等核心事件,丢弃策略必须显式记录;
- 延长榜单快照发布时间,继续服务上一版本;
- 对异常热点 topic 单独隔离;
- 若核心数据源不完整,停止发布新榜并告警。
降级优先牺牲实时性和次要特征,不优先牺牲可恢复的核心行为事实。
十四、安全、隐私与运营审计
14.1 最小化行为数据
热搜计算只使用完成聚合所需的主体键、事件类型、话题和时间。主体键应脱敏或令牌化;原始查询文本、设备标识和地域数据按权限、用途和保留期限隔离。
14.2 控制台权限
| 操作 | 建议权限和控制 |
|---|---|
| 查看自然热度和解释 | 榜单运营、数据分析只读 |
| 修改算分规则 | 双人审核、版本化、灰度和回滚 |
| 置顶或保底 | 独立运营权限,记录有效期和原因 |
| 紧急下榜 | 内容安全高优先级,事后复核 |
| 解除屏蔽 | 与下榜操作分权或二次审核 |
| 历史行为查询 | 严格审计和数据最小化 |
14.3 审计不可混入自然分
所有人工动作保存前后状态、操作者、审批、原因、工单和有效时间。展示位置可以受运营策略影响,但内部必须能同时看到自然排名和最终展示排名。
十五、可观测、对账与故障排查
15.1 技术指标
| 层级 | 指标 |
|---|---|
| 事件源 | 各类型 EPS、失败、Outbox 未投递数、schema 版本 |
| 消息流 | 分区积压、最老事件年龄、热点分区、重复率 |
| 归一 | 未识别词比例、低置信候选、mapping_version、mapping activation epoch、归一前 barrier 向量、旧 epoch 重路由与各 scope cutover 状态 |
| 窗口 | 各 scope/source/partition 的 watermark、idle 与恢复、head/不可变版本、aggregation generation、迟到比例、检查点耗时 |
| 算分 | COLLECTING/BUILDING/SEALED run 数与年龄、预期分区未完成数、分区 Top K 摘要、run 输入集合未终态数、候选集合不闭合、异常分数、score_version、TopK 变化率和 CAS 拒绝数 |
| 风控 | 各 scope 可疑流量、risk epoch、global risk epoch、到期任务延迟、边缘 epoch 落后、purge 失败、误判复核 |
| 发布 | 各 scope 快照延迟、PREPARED/READY 滞留、三类 activation epoch、publication_epoch、context/pointer CAS 冲突 |
| 查询 | QPS、缓存命中、P95/P99、旧快照返回比例 |
15.2 业务健康指标
- 榜单话题的有效主体分布;
- 新话题从达到门槛到上榜的延迟;
- 榜单平均停留时长和名次抖动;
- 自然榜与展示榜的差异比例;
- 用户点击后的负反馈和快速返回;
- 下榜申诉、误杀和人工修正数量;
- 相似话题重复占榜比例;
- 数据源缺失时仍发布的新快照数。
15.3 三层对账
每次对账先冻结比较上下文:
scope_id + input_watermark + window_config_version
+ expected_partition_set_digest + partition_result_set_digest
+ input_manifest_set_digest
+ rule_version + rule_activation_epoch
+ mapping_version + mapping_activation_epoch
+ risk_epoch + risk_model_version + global_risk_epoch
+ calculation_run_id
然后再比较三层:
- 该 scope 对应的各源事件位点与原始归档是否连续,idle 和恢复区间是否已解释;
- 同一水位下抽样话题的原始事件、有效去重和完整键窗口聚合是否可重算;
- 同一预期候选分区集合、分区 cut、输入清单集合和 activation 上下文下的候选、PUBLISHED 快照、数据库权威指针、Redis 指针和 CDN 响应是否一致。
不固定上述上下文时,数量差异可能只是比较了不同水位或规则,不能据此判断数据损坏。对账报告要保存抽样范围、源位点摘要、差异明细和结论状态。
15.4 故障排查示例
现象:某个普通话题在一分钟内从榜外冲到第一。
- 固定异常响应的
scope_id和publication_epoch,检查快照、输入清单、三类 activation epoch 和运营调整; - 展开自然分,看增长来自哪个事件类型和窗口;
- 检查去重主体、账号年龄、设备/IP 聚集和相似内容;
- 对比原始事件速率与上游业务事实,排除重复投递;
- 检查话题是否误合并了另一个热门事件;
- 风险明确时走审计下榜,不直接删除数据;
- 修复去重、归一或消费位点后,在隔离环境回放验证;
- 生成修复快照,并记录原因和影响范围。
十六、测试、故障演练与验证证据
16.1 正确性测试
以下都是本设计的验收计划,不代表已经执行。只有保存输入数据、配置、断言结果和报告位置后,才能把状态改为 PASS 或 FAIL。
| ID | 场景 | 通过条件 | 状态 |
|---|---|---|---|
| C01 | 同一 event_id 多次投递 | 完整键窗口的事件数、有效值和 aggregate_version 按幂等策略收敛 | NOT_RUN |
| C02 | 同一 actor 短时间重复搜索 | 按主体规则合并或降权,原始值与基础有效值仍可解释 | NOT_RUN |
| C03 | 同一事件存在多个别名 | 当前 mapping_version 下只映射到一个主话题并只占一个榜位 | NOT_RUN |
| C04 | global、channel、region、language 四类 scope 同时计算 | 同一 topic 在各 scope_id 的窗口、候选、快照和风险动作互不覆盖 | NOT_RUN |
| C05 | 相同 snapshot 参数请求不同 scope | Redis Key、ETag、CDN cache key 和响应均不跨 scope 命中 | NOT_RUN |
| C06 | 热点局部结果具有相同 topic 但不同 scope/window/metric | 只按完整聚合键合并,未出现跨 scope、窗口或指标串值 | NOT_RUN |
| C07 | 话题合并期间源话题和主话题持续收到并发事件 | 影子状态包含窗口与 actor 去重状态;barrier 前后核对无重复、无漏算,历史快照不变 | NOT_RUN |
| C08 | barrier 前后混入旧 event_time 的迟到事件 | 事件只按各分区 cutover_next_offset 判归属;barrier 后的迟到事件不返回旧 merge run | NOT_RUN |
| C09 | 低流量分区进入 idle 后恢复并追赶 | watermark 不倒退,积压按迟到策略处理,通过位点检查后才重新纳入 | NOT_RUN |
| C10 | 去重状态到期后注入超出自动重放范围的历史事件 | 正式聚合不增加,事件只进入隔离回放或审计 | NOT_RUN |
| C11 | 分别向 CLOSED 和 FINAL 窗口发送迟到事件 | CLOSED 提升聚合版本并只影响后续快照;FINAL 正式结果不变 | NOT_RUN |
| C12 | 风控响应晚于新窗口、风险动作或规则激活 | 输入清单或任一 activation epoch 不符时 CAS 失败,延迟结果不能覆盖新候选 | NOT_RUN |
| C13 | 旧规则任务、旧发布器和未授权 replay run 同时抢占指针 | scope context version、三类 epoch、水位和 run 授权任一不符都被拒绝 | NOT_RUN |
| C14 | 快照已 READY 但指针 CAS 未成功,此时 Redis 丢失 | 只恢复数据库权威指针指向的 PUBLISHED 快照,不返回该 READY 快照 | NOT_RUN |
| C15 | 在频道 A 屏蔽并解除一个话题 | A 的边缘与查询过滤生效且每次推进 risk epoch,频道 B 不受局部动作影响 | NOT_RUN |
| C16 | 多个话题分数完全相同并重复计算 | 按稳定次级字段得到一致顺序 | NOT_RUN |
| C17 | 候选榜缓存丢失 | 固定 scope context,从窗口聚合生成输入清单并重建;上下文变化时整批废弃 | NOT_RUN |
| C18 | 已发布快照缓存丢失 | 按数据库权威指针和 publication_epoch 恢复同一不可变快照 | NOT_RUN |
| C19 | 三层对账故意混入不同 scope 或版本数据 | 对账任务拒绝比较;固定同一上下文后才输出差异 | NOT_RUN |
| C20 | 多 scope 合并在 barrier 前、部分 scope CAS 后和全部完成前分别宕机 | 每条 scope cutover 独立恢复,已切 scope 不回退,未切 scope 不抢跑,父 run 不冒充全局完成 | NOT_RUN |
| C21 | 算分冻结 1 分钟、5 分钟和 1 小时多指标后,窗口 head 继续推进 | 旧 run 仍按不可变版本和值重现;新 head 只触发新 run,不让旧 run 混读或无故失效 | NOT_RUN |
| C22 | rule-v2 升级到 v3 后又回滚 v2,旧 v2 任务此时恢复 | 回滚分配更高 rule activation epoch,旧 v2 任务因 epoch 不同被拒绝 | NOT_RUN |
| C23 | 全局风险动作与频道局部动作在发布过程中并发 | global/scope risk epoch 任一变化都使旧快照 CAS 失败,边缘和查询过滤立即按新动作收敛 | NOT_RUN |
| C24 | 一个计算 run 只写完一半候选就宕机,发布器同时扫描 | 未终态输入使 run 无法 SEALED;恢复后精确集合闭合或新建 run,快照不混用两轮候选 | NOT_RUN |
| C25 | 候选 worker 提交最后一项时封口事务并发执行 | 二者锁同一 BUILDING run;封口要么看到完整终态集合成功,要么失败重试,SEALED 后无迟到写 | NOT_RUN |
| C26 | 删除一个候选并用另一个重复条目维持相同数量 | 输入集合摘要、topic 集合和 disposition 对不上,封口失败,不能只靠 count 通过 | NOT_RUN |
| C27 | barrier 后旧 normalizer 才输出旧 epoch 结果 | 聚合器拒绝旧 epoch,按归一前 source offset 幂等重路由,新旧 topic 不双算 | NOT_RUN |
| C28 | CDN 保持热缓存时执行全局紧急下榜,并故意让 purge 失败 | 所有受影响 scope 仍经边缘 denylist 过滤;边缘 epoch 落后时绕过缓存,旧完整响应不直出 | NOT_RUN |
| C29 | 局部与全局风险动作到期任务延迟、重复执行 | 未提交 EXPIRED 前保持屏蔽;幂等到期事务只推进一次相应 epoch,缓存有效期不跨到期边界 | NOT_RUN |
| C30 | 新 aggregation generation 切换成功后暂停的旧 worker 恢复写 head | 旧 worker 因 active generation fence 失败,不能用更大数字版本推进新 head;两代版本主键与 manifest 引用不冲突 | NOT_RUN |
| C31 | 一个预期候选分区完全不提交 Top K,其余分区和 manifest 内部都闭合 | run 停在 COLLECTING,缺失分区不能被剩余集合摘要掩盖,发布器继续服务旧快照 | NOT_RUN |
| C32 | 分区 Top K 在 run 已进入 BUILDING 后迟到或尝试改写摘要 | 不可变 partition cut 拒绝迟到写;新结果只能进入新 run,既有 manifest 集合不变化 | NOT_RUN |
| C33 | 快照复制了正确的两个集合摘要,但故意少写 input_manifest_count 或 candidate_count | 发布前校验拒绝快照,两个 count 必须分别等于 SEALED run 的冻结输入数和候选数,不能由 Top N 榜项数替代 | NOT_RUN |
16.2 算法离线评估
用固定 scope、输入水位、候选输入清单、映射 activation epoch 和风险 epoch 的历史事件回放比较新旧规则:
| ID | 评估项 | 记录指标 | 状态 |
|---|---|---|---|
| A01 | Top N 稳定性 | 重合率、名次变化和退出数量 | NOT_RUN |
| A02 | 突发事件发现 | 达到门槛至上榜的延迟 | NOT_RUN |
| A03 | 老话题霸榜 | 在榜时长、基线修正前后排名 | NOT_RUN |
| A04 | 机器人样本 | 异常样本入榜率和风险调整幅度 | NOT_RUN |
| A05 | 风险准确性 | 风险话题拦截率和正常话题误伤率 | NOT_RUN |
| A06 | 榜单抖动 | 快照 Top N 变化率和平均停留时间 | NOT_RUN |
| A07 | scope 偏差 | 不同频道、地域和语言下的召回与分布差异 | NOT_RUN |
离线结果只能证明历史样本表现,不能替代线上灰度和数据质量监控。
16.3 容量测试
| ID | 压测场景 | 重点观测 | 状态 |
|---|---|---|---|
| P01 | 均匀长尾 topic | EPS、CPU、网络和状态增长 | NOT_RUN |
| P02 | 单一突发 topic 占大比例事件 | 热点分区延迟、局部合并积压和公平性 | NOT_RUN |
| P03 | 大量重复 actor 与 event | 去重吞吐、状态内存和误计数 | NOT_RUN |
| P04 | global/channel/region/language 多 scope 扇出 | 展开倍率、分区倾斜和 scope 隔离 | NOT_RUN |
| P05 | 消息分区重平衡 | watermark 停顿、重复消费和恢复时长 | NOT_RUN |
| P06 | 窗口及去重状态接近容量上限 | 状态大小、GC、检查点时长和清理速度 | NOT_RUN |
| P07 | 候选 Top K 快速变化 | 候选写入、过期 CAS 拒绝和发布延迟 | NOT_RUN |
| P08 | 快照发布与查询高峰同时发生 | 数据库 CAS、Redis、CDN 和查询 P99 | NOT_RUN |
| P09 | 一个计算或查询节点退出 | 剩余容量、积压和错误率 | NOT_RUN |
| P10 | 从检查点和最长自动重放范围追赶 | 追赶速度、水位恢复和重复率 | NOT_RUN |
每项报告都应注明事件分布、事件字节、scope 数、窗口配置、实例规格、持续时间、分位延迟、积压、状态大小、检查点时间和恢复时间。当前全部为容量验证计划,不能据此声称达到某个吞吐量。
16.4 故障演练
| ID | 注入故障 | 预期行为 | 状态 |
|---|---|---|---|
| F01 | 某一必需事件源停止发送 | 标记数据不完整,冻结发布或切换显式降级规则 | NOT_RUN |
| F02 | 分区进入 idle 后恢复且携带大量旧事件 | 不倒退 watermark,旧事件按迟到边界分流 | NOT_RUN |
| F03 | Kafka 分区积压、重平衡或重复回放 | 位点连续,幂等聚合收敛,追赶时间可观测 | NOT_RUN |
| F04 | 话题归一服务超时 | 原事件可重试,低置信或未映射事件不污染候选 | NOT_RUN |
| F05 | 风控服务不可用或响应严重延迟 | 按版本化策略降级,旧响应无法越过 scope risk epoch 覆盖新候选 | NOT_RUN |
| F06 | 流处理检查点损坏 | 从上一一致检查点恢复,完整聚合键和位点对账通过 | NOT_RUN |
| F07 | Redis 候选榜丢失 | 从窗口事实重建到隔离 namespace 后再替换 | NOT_RUN |
| F08 | Redis 快照和 current pointer 丢失 | 只从数据库权威发布指针恢复 PUBLISHED 快照 | NOT_RUN |
| F09 | 快照 READY 后发布器在指针 CAS 前宕机 | 继续服务旧 PUBLISHED 快照,READY 不被查询 | NOT_RUN |
| F10 | 旧发布器在新榜发布后恢复 | publication epoch、scope context version 和三类 activation epoch 拒绝旧写 | NOT_RUN |
| F11 | CDN 热缓存未及时失效或 purge 失败 | 命中后边缘过滤仍按最新 scope/global denylist 下榜;epoch 落后就绕过缓存回源 | NOT_RUN |
| F12 | 规则灰度异常后回滚 | 旧 rule_version 以更高 activation epoch 重新激活,不复用旧任务或旧 epoch | NOT_RUN |
| F13 | 未授权历史 replay 尝试发布 | run 授权校验失败,正式候选和指针不变 | NOT_RUN |
| F14 | 风险到期调度暂停后恢复并重复投递 | 动作在 EXPIRED 事务前保持有效,恢复后只推进一次 risk epoch 并触发新快照 | NOT_RUN |
| F15 | 候选生产阶段一个预期分区卡死或整批结果丢失 | run 不进入 BUILDING;只有新建携审批降级策略和明确缺失集合的 run 才可能发布 | NOT_RUN |
十七、设计取舍与演进
17.1 定时批处理什么时候够用
社区规模小、更新容忍分钟级时,可以每分钟从增量行为表聚合候选并生成快照,结构简单且容易审计。只有事件量、实时性或窗口复杂度被验证为瓶颈时,再引入完整流处理平台。
17.2 Redis ZSET 什么时候够用
如果只有单一指标、允许数据丢失后重建、没有复杂去重和规则审计,ZSET 可以做候选榜。但仍需要保存行为或窗口事实、发布快照和风险操作。ZSET 是排序索引,不是业务全貌。
17.3 精确去重和近似去重
| 方案 | 优点 | 代价和边界 |
|---|---|---|
| 精确 actor 集合 | 可列明细、误差小 | 热点窗口内存大,隐私和保留压力高 |
| HyperLogLog 等近似结构 | 可合并、内存稳定 | 有估算误差,不能列出主体 |
| 分层策略 | 长尾近似、候选精算 | 实现和口径更复杂,但更平衡 |
热搜排序本身允许小误差时,可使用近似基数;涉及风控调查或申诉时,从受控原始证据抽样核查,不能把 HLL 当明细库。
17.4 演进触发条件
- 批处理无法达到已定义更新时效时引入流式窗口;
- 单 topic 打满分区时引入局部预聚合和动态拆分;
- 规则频繁变更且无法解释时强化规则平台、影子和审计;
- 机器人损失或误伤达到业务阈值时升级群体风控;
- 全站榜、频道榜和地域榜的计算逻辑明显分化时拆分 scope 资源;
- 原始事件存储成本过高时调整分层保留,而不是删除所有恢复证据。
十八、面试复盘与职责边界
18.1 三分钟主回答
我不会把热搜直接等同 Redis ZSET。先统一搜索、发帖、评论和转发等行为事件,经契约校验、事件去重和话题归一后,按 topic 和事件时间做多尺度窗口聚合。指标同时看有效行为、去重主体、相对历史基线的增速、内容质量和时间衰减,算分规则必须版本化。机器人行为通过主体去重、群体异常和风险因子校正,原始值与校正值分开保存。
一次候选算分会先固定候选生产的预期分区和 source cut,收齐每个分区的 Top K 摘要后,再把多窗口、多指标及各自 generation + aggregate_version 固化成输入清单;run 随后冻结完整 manifest 集合,而不是只记一个最大版本或候选数量。发布器固定 scope 上下文和水位读取一个 SEALED run,先完整写不可变榜单快照,再条件锁定但不修改 scope context,并 CAS 数据库权威指针,客户端始终读一张完整榜。规则、话题映射和风控各有单调 activation epoch;即使回滚到同一份旧规则,也分配新 epoch,旧任务不能借同名版本复活。话题合并使用归一前稳定日志的分区 barrier 位点按 scope 划交接所有权,event time 只决定窗口迟到策略。
迟到事件在允许范围内修正后续快照,超范围进入离线审计;Redis 丢失可由窗口聚合、输入清单和已发布快照恢复。规则升级先影子回放、比较差异和灰度,异常时用更高 activation epoch 重新激活旧内容。容量测试特别覆盖单一热点话题、重复事件、分区积压和检查点恢复。
18.2 高频追问
- 为什么不用数据库
order by count desc?用户请求时全量聚合代价高,也无法处理实时窗口、基线、去重和风险校正。 - 为什么单纯 ZSET 不够?它只有当前分数,没有原始证据、规则版本、历史快照和恢复链路。
- 老话题为什么不会一直第一?使用相对基线、增长速度、时间衰减和在榜疲劳,而不是只看累计总量。
- 同一用户刷一万次怎么算?短窗口主体去重或降权,结合设备、IP 和群体行为识别,不只看账号。
- 热点 topic 把单分区打满怎么办?局部预聚合、actor 稳定分片、可合并去重结构和热点资源隔离。
- 消息乱序怎么办?按 event_time 进窗口,用 watermark 和允许迟到范围收敛。
- 迟到一天的事件要改昨天榜单吗?通常不改已发布产品事实,进入离线评估或审计版本。
- 规则改动怎么发布?新版本固定同一输入清单做影子比较,灰度后提升 activation epoch;回滚旧内容也使用新 epoch。
- 风控下榜会不会丢失真实热度?不会,风险操作和自然分分开,原始指标保留;scope risk epoch 会阻止旧快照越过新动作发布。
- Redis 挂了怎么办?候选榜从窗口聚合重算,正式榜从不可变快照恢复。
- 怎样避免榜单频繁抖动?多尺度窗口、进入退出不同阈值、连续低分才退出和稳定同分排序。
- 怎样证明结果可信?保存源位点、窗口版本向量、输入清单、三类 activation epoch、快照和风险动作,固定上下文后抽样回放对账。
- 话题合并为什么不能只按 event time 或归一后 offset 切换?event time 会迟到,topic 又会改变下游分区;所有权必须由归一前稳定日志的 barrier offset 唯一划分。
- 多窗口算分为什么不能只存一个 aggregate version?各聚合行独立推进,清单还要引用包含实际值和 source cut 的不可变版本,才能在 head 前进后重现。
- run 有 candidate count 为什么还可能少写?少一项再重复另一项时数量可以相同;既要先证明所有预期候选分区都在同一 source cut 提交 Top K,再冻结输入 manifest 精确集合,并让每项都落成候选或带原因的拒绝。
- 发布时为什么不推进 scope context version?发布只消费既定控制上下文;它条件锁定版本并 CAS 指针,只有规则、映射和风险控制事务负责推进 context。
- CDN purge 失败时紧急下榜怎么保证?命中后还要经过带 risk epoch 的边缘 denylist;边缘版本落后就绕过缓存,无法做边缘过滤时只能缩短 TTL 并在紧急状态绕过 CDN。
- 恢复代次切换后旧 worker 为什么不能靠更大的 aggregate_version 抢回 head?版本只在 generation 内比较,head CAS 同时校验 active generation 和当前版本,manifest 也引用完整三元组。
18.3 可选择的职责主线
如果项目真实存在,可把个人职责拆成互不冒领的方向:
- 事件与流计算:负责契约、窗口、watermark、检查点和容量验证;
- 算分与发布:参与规则版本、候选榜、快照原子发布和缓存恢复;
- 风控与运营:协作主体去重、风险动作、控制台审计和紧急下榜。
只实现过 Redis 排行榜,就应明确讲到“当前方案”和可演进边界。没有真实事件规模、回放报告和监控证据,不能声称支撑过某个超大平台的生产热搜。
18.4 自检问题
- 热度和累计行为总量有什么区别?
- 为什么需要稳定
topic_id? - 新话题如何获得机会又不让垃圾词上榜?
- 去重主体为什么可能用近似算法?
- 同一个 actor 落到多个局部分区时怎样合并?
- watermark 和处理时间分别解决什么?
- 哪些迟到事件会修正在线结果?
- 规则版本如何与快照关联?
- 风险分与自然分为什么要分开?
- 快照发布为什么先写内容再切指针?
- 旧发布任务如何避免覆盖新快照?
- 某数据源中断时为什么不能照常套原公式?
- 候选榜丢失和正式榜丢失分别从哪里恢复?
- 哪些指标证明榜单没有被机器人主导?
- 什么时候批处理比流处理更合适?
- 话题合并时,barrier offset 和 effective watermark 分别解决什么?
- 一次候选算分读取多个窗口时,怎样证明没有混入不同水位的数据?
- 回滚到旧 rule_version 时,为什么仍要分配更高 activation epoch?
- 一个话题的风险动作如何让正在生成的旧 scope 快照失效?
- 窗口 head 前进后,旧 run 依靠什么重现当时的输入值?
SEALED怎样证明输入集合没有少处理一项,且封口后没有 worker 迟到写?- 旧 normalizer 结果跨过合并 barrier 才到达时,怎样拒绝并重路由?
- 风险动作到期任务延迟时,为什么不能只按墙上时间自动解除?
- 一个候选分区整批未上报时,run 为什么不能仅凭已收到的 manifest 集合封口?
- aggregation generation 切换后,旧 worker 和旧 manifest 分别由什么 fence 阻止污染新代次?
总结
实时热搜榜的核心不是“给话题加分”,而是建立一条可解释、可恢复的实时决策链:
- 用可靠行为事件和稳定话题标识解决输入事实;
- 用多尺度窗口、主体去重、历史基线和时间衰减表达“此刻变热”;
- 用归一前稳定日志的分区 barrier、按 scope 的交接记录和映射 activation epoch 让话题合并只有一条事件负责链路;
- 用带 generation 的不可变窗口版本、候选分区 cut、输入清单和 run 级精确集合证明一次算分真正收齐、读取并处理了哪些输入;
- 用 scope 风控纪元、显式到期事务和边缘 denylist 控制刷量、缓存与内容安全;
- 用可封口候选 run、数据库权威指针和不可变快照分离持续计算与稳定发布;
- 用事件位点、窗口聚合、快照和隔离回放处理故障与争议。
Redis ZSET 可以是候选排序组件,但原始事件、窗口指标和已发布快照才让系统在重复、迟到、规则变化、热点和缓存故障后仍然说得清、恢复得了。