返回资源中心

实时热搜榜设计:窗口聚合、抗刷与快照发布

技术文章从话题归一到实时窗口聚合,设计可解释、抗刷、可恢复的热搜榜,覆盖迟到事件、候选集合、风险纪元、快照发布和对账。小红学堂更新于 2026年8月29日41 次阅读

内容平台的实时热搜榜如何设计?

阅读说明与事实边界

“热搜榜怎么做”经常被答成一句话:把话题和分数放进 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. 当前窗口相对历史基线的增长速度;
  3. 参与用户的去重数量,而不是单纯点击次数;
  4. 多种行为之间的质量差异;
  5. 时间衰减,让旧事件自然退出;
  6. 风险、内容安全和人工运营约束。

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_epocheffective_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 创建时为每个预期候选分区插入一行 PENDINGcandidate_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_digesttopic_id + input_manifest_id + manifest_digest 稳定排序计算;每个条目最终只能进入 CANDIDATE 或带原因的 REJECTED,从而能证明没有漏处理某个输入;
  • 候选以 scope_id + calculation_run_id + topic_id 隔离,写候选与确认父 run 仍为 BUILDING 放在同一事务;数据库拒绝对 SEALED run 新增、更新或删除候选。发布器只读取一个已封口的正式 run;影子规则和回放任务另放独立 namespace,并携带输入清单与三类 activation epoch;
  • 快照项以 snapshot_id + rank_no 唯一;
  • 同一快照内 topic_id 唯一,防止合并话题重复上榜;
  • 风险操作按 scope_id + topic_id + effective_from 索引;跨所有榜的动作使用明确的全局 scope,再由查询规则与局部 scope 共同求值;
  • RANKING_SCOPE_CONTEXTscope_id 唯一,保存当前激活的规则、窗口配置、映射、scope 风控纪元及单调控制版本;规则回滚也只能继续增加 activation epoch;
  • RANKING_PUBLISH_POINTERscope_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-CNCHANNEL|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 紧急下榜

内容安全确认需要立即下榜时:

  1. 在控制事务中写入带 scope_id 和递增 risk_action_version 的风险操作,同时推进该 scope 的 risk_epoch
  2. 把带 epoch 的高优先级 denylist 推送到 CDN 命中之后仍会执行的边缘过滤层,并同步到查询层;边缘尚未确认最新 epoch 时绕过榜单缓存回源,不等待下一张榜;
  3. 候选按新的风险 epoch 重算,发布器生成不包含该话题的新快照;
  4. 清理所有受影响 scope 的 CDN 对象;全局动作展开受影响 scope 清单,purge 失败时由边缘 denylist 继续过滤,缓存最大安全 TTL 作为最后边界;
  5. 记录操作人、规则、原因、动作版本、风险 epoch 和生效时间;
  6. 正常计算继续保留自然分,用于复核,不自动对外展示。

紧急通道必须可审计且有显式到期事务,避免一次临时屏蔽永久遗留。解除屏蔽或到期同样生成新动作版本、推进 epoch、触发候选重算和新快照,不能直接删掉 Redis 屏蔽 member;全局动作和局部动作的并集仍在当前 scope 下求值,避免频道 A 的操作误伤频道 B。若平台没有 CDN 命中后的可编程过滤能力,安全口径只能选择短 TTL 并在紧急状态下绕过 CDN,不能声称“查询层一定即时拦截”。

9.4 运营干预

置顶、保底曝光、活动话题位和自然热搜最好分成不同产品位。如果必须混排,快照项同时保存自然分、调整类型和最终位置。运营不能直接把数据库的自然分改成一个更大数字。

十、候选榜与快照发布

10.1 为什么分两阶段

流处理不断更新候选分数,而对外榜单需要稳定版本。候选不能只按 scope_id + topic_id 原地覆盖,否则本轮只更新了一半时,发布器可能把上一轮和本轮拼成一张榜。每轮正式计算使用独立 calculation_run_id

  1. 创建 run 时先冻结候选生产阶段的 source_offset_vector_ref 和预期分区集合,按 candidate_partition_idRUN_PARTITION_CUT(PENDING),保存每个分区可验证的 expected_source_cut_refexpected_partition_set_digest(candidate_partition_id, expected_source_cut_ref) 稳定排序计算,不能只摘要分区 ID;run 此时是 COLLECTING
  2. 每个分区只计算自己固定 cut 以内的局部 Top K,把不可变的 topk_object_ref + topk_item_count + topk_digest 一次性提交后,才条件推进为 COMPLETED。重复提交只有内容完全相同才按幂等成功,整个分区未上报不能靠其他分区的完整结果掩盖;
  3. 协调器确认 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,再 CAS COLLECTING -> BUILDING
  4. BUILDING 阶段只允许处理这个冻结集合。每个输入条目必须落成 CANDIDATE 或带原因的 REJECTED,不能靠“最后大约有 N 条”判断完成;
  5. worker 写候选、写 result digest 和把条目标终态时,在同一事务内锁定父 run 并确认仍为 BUILDING
  6. 封口事务锁定 run,验证期望输入集合一个不少、一个不多,所有条目均为终态,候选 topic 集合与 CANDIDATE 条目精确相等,再写 candidate_count + candidate_set_digest 并 CAS 为 SEALED
  7. BUILDING 后分区 cut 和输入集合不可修改;SEALED 后数据库再拒绝候选的新增、修改或删除,迟到修正只能生成新 run;
  8. 发布器只从一个 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,竞争失败、版本过期或校验失败进入 REJECTEDREADY 只说明内容完整,不代表产品已经发布;数据库必须在同一事务中完成 RANKING_PUBLISH_POINTER 的 CAS 和快照状态改为 PUBLISHED,任一步失败都回滚。写快照内容和切换权威指针必须分开,客户端不会读到半张榜,也不会在 Redis 丢失时误选一个从未发布成功的 READY 快照。

每个 scope_id 只有一条数据库权威发布指针和一条 scope 控制上下文。前者回答“用户当前读哪张已发布榜”,后者回答“现在允许哪个规则、映射和风控上下文发布”。两者不能混成 Redis 里的一条可随意覆盖记录。基线把 context、pointer 和 snapshot 状态放在同一事务数据库中:发布事务先条件锁定 context 行并验证版本,不增加 context version,再 CAS 指针并发布快照;规则、映射和风险事务也锁同一 context 行并负责推进它的版本。若它们被拆到不同数据库,就必须另设单一发布仲裁服务,不能把两次本地 CAS 冒充原子提交。发布事务至少同时验证:

  • 指针的 scope_id 与快照 scope 完全一致;
  • 读取时的 versionpublication_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&regionId=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 背压与降级顺序

  1. 暂停非关键解释特征和低价值曝光事件;
  2. 扩大局部聚合批次,降低候选榜更新频率;
  3. 保留搜索、发帖等核心事件,丢弃策略必须显式记录;
  4. 延长榜单快照发布时间,继续服务上一版本;
  5. 对异常热点 topic 单独隔离;
  6. 若核心数据源不完整,停止发布新榜并告警。

降级优先牺牲实时性和次要特征,不优先牺牲可恢复的核心行为事实。

十四、安全、隐私与运营审计

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

然后再比较三层:

  1. 该 scope 对应的各源事件位点与原始归档是否连续,idle 和恢复区间是否已解释;
  2. 同一水位下抽样话题的原始事件、有效去重和完整键窗口聚合是否可重算;
  3. 同一预期候选分区集合、分区 cut、输入清单集合和 activation 上下文下的候选、PUBLISHED 快照、数据库权威指针、Redis 指针和 CDN 响应是否一致。

不固定上述上下文时,数量差异可能只是比较了不同水位或规则,不能据此判断数据损坏。对账报告要保存抽样范围、源位点摘要、差异明细和结论状态。

15.4 故障排查示例

现象:某个普通话题在一分钟内从榜外冲到第一。

  1. 固定异常响应的 scope_idpublication_epoch,检查快照、输入清单、三类 activation epoch 和运营调整;
  2. 展开自然分,看增长来自哪个事件类型和窗口;
  3. 检查去重主体、账号年龄、设备/IP 聚集和相似内容;
  4. 对比原始事件速率与上游业务事实,排除重复投递;
  5. 检查话题是否误合并了另一个热门事件;
  6. 风险明确时走审计下榜,不直接删除数据;
  7. 修复去重、归一或消费位点后,在隔离环境回放验证;
  8. 生成修复快照,并记录原因和影响范围。

十六、测试、故障演练与验证证据

16.1 正确性测试

以下都是本设计的验收计划,不代表已经执行。只有保存输入数据、配置、断言结果和报告位置后,才能把状态改为 PASSFAIL

ID场景通过条件状态
C01同一 event_id 多次投递完整键窗口的事件数、有效值和 aggregate_version 按幂等策略收敛NOT_RUN
C02同一 actor 短时间重复搜索按主体规则合并或降权,原始值与基础有效值仍可解释NOT_RUN
C03同一事件存在多个别名当前 mapping_version 下只映射到一个主话题并只占一个榜位NOT_RUN
C04global、channel、region、language 四类 scope 同时计算同一 topic 在各 scope_id 的窗口、候选、快照和风险动作互不覆盖NOT_RUN
C05相同 snapshot 参数请求不同 scopeRedis Key、ETag、CDN cache key 和响应均不跨 scope 命中NOT_RUN
C06热点局部结果具有相同 topic 但不同 scope/window/metric只按完整聚合键合并,未出现跨 scope、窗口或指标串值NOT_RUN
C07话题合并期间源话题和主话题持续收到并发事件影子状态包含窗口与 actor 去重状态;barrier 前后核对无重复、无漏算,历史快照不变NOT_RUN
C08barrier 前后混入旧 event_time 的迟到事件事件只按各分区 cutover_next_offset 判归属;barrier 后的迟到事件不返回旧 merge runNOT_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
C22rule-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
C27barrier 后旧 normalizer 才输出旧 epoch 结果聚合器拒绝旧 epoch,按归一前 source offset 幂等重路由,新旧 topic 不双算NOT_RUN
C28CDN 保持热缓存时执行全局紧急下榜,并故意让 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_countcandidate_count发布前校验拒绝快照,两个 count 必须分别等于 SEALED run 的冻结输入数和候选数,不能由 Top N 榜项数替代NOT_RUN

16.2 算法离线评估

用固定 scope、输入水位、候选输入清单、映射 activation epoch 和风险 epoch 的历史事件回放比较新旧规则:

ID评估项记录指标状态
A01Top N 稳定性重合率、名次变化和退出数量NOT_RUN
A02突发事件发现达到门槛至上榜的延迟NOT_RUN
A03老话题霸榜在榜时长、基线修正前后排名NOT_RUN
A04机器人样本异常样本入榜率和风险调整幅度NOT_RUN
A05风险准确性风险话题拦截率和正常话题误伤率NOT_RUN
A06榜单抖动快照 Top N 变化率和平均停留时间NOT_RUN
A07scope 偏差不同频道、地域和语言下的召回与分布差异NOT_RUN

离线结果只能证明历史样本表现,不能替代线上灰度和数据质量监控。

16.3 容量测试

ID压测场景重点观测状态
P01均匀长尾 topicEPS、CPU、网络和状态增长NOT_RUN
P02单一突发 topic 占大比例事件热点分区延迟、局部合并积压和公平性NOT_RUN
P03大量重复 actor 与 event去重吞吐、状态内存和误计数NOT_RUN
P04global/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 和查询 P99NOT_RUN
P09一个计算或查询节点退出剩余容量、积压和错误率NOT_RUN
P10从检查点和最长自动重放范围追赶追赶速度、水位恢复和重复率NOT_RUN

每项报告都应注明事件分布、事件字节、scope 数、窗口配置、实例规格、持续时间、分位延迟、积压、状态大小、检查点时间和恢复时间。当前全部为容量验证计划,不能据此声称达到某个吞吐量。

16.4 故障演练

ID注入故障预期行为状态
F01某一必需事件源停止发送标记数据不完整,冻结发布或切换显式降级规则NOT_RUN
F02分区进入 idle 后恢复且携带大量旧事件不倒退 watermark,旧事件按迟到边界分流NOT_RUN
F03Kafka 分区积压、重平衡或重复回放位点连续,幂等聚合收敛,追赶时间可观测NOT_RUN
F04话题归一服务超时原事件可重试,低置信或未映射事件不污染候选NOT_RUN
F05风控服务不可用或响应严重延迟按版本化策略降级,旧响应无法越过 scope risk epoch 覆盖新候选NOT_RUN
F06流处理检查点损坏从上一一致检查点恢复,完整聚合键和位点对账通过NOT_RUN
F07Redis 候选榜丢失从窗口事实重建到隔离 namespace 后再替换NOT_RUN
F08Redis 快照和 current pointer 丢失只从数据库权威发布指针恢复 PUBLISHED 快照NOT_RUN
F09快照 READY 后发布器在指针 CAS 前宕机继续服务旧 PUBLISHED 快照,READY 不被查询NOT_RUN
F10旧发布器在新榜发布后恢复publication epoch、scope context version 和三类 activation epoch 拒绝旧写NOT_RUN
F11CDN 热缓存未及时失效或 purge 失败命中后边缘过滤仍按最新 scope/global denylist 下榜;epoch 落后就绕过缓存回源NOT_RUN
F12规则灰度异常后回滚旧 rule_version 以更高 activation epoch 重新激活,不复用旧任务或旧 epochNOT_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 高频追问

  1. 为什么不用数据库 order by count desc?用户请求时全量聚合代价高,也无法处理实时窗口、基线、去重和风险校正。
  2. 为什么单纯 ZSET 不够?它只有当前分数,没有原始证据、规则版本、历史快照和恢复链路。
  3. 老话题为什么不会一直第一?使用相对基线、增长速度、时间衰减和在榜疲劳,而不是只看累计总量。
  4. 同一用户刷一万次怎么算?短窗口主体去重或降权,结合设备、IP 和群体行为识别,不只看账号。
  5. 热点 topic 把单分区打满怎么办?局部预聚合、actor 稳定分片、可合并去重结构和热点资源隔离。
  6. 消息乱序怎么办?按 event_time 进窗口,用 watermark 和允许迟到范围收敛。
  7. 迟到一天的事件要改昨天榜单吗?通常不改已发布产品事实,进入离线评估或审计版本。
  8. 规则改动怎么发布?新版本固定同一输入清单做影子比较,灰度后提升 activation epoch;回滚旧内容也使用新 epoch。
  9. 风控下榜会不会丢失真实热度?不会,风险操作和自然分分开,原始指标保留;scope risk epoch 会阻止旧快照越过新动作发布。
  10. Redis 挂了怎么办?候选榜从窗口聚合重算,正式榜从不可变快照恢复。
  11. 怎样避免榜单频繁抖动?多尺度窗口、进入退出不同阈值、连续低分才退出和稳定同分排序。
  12. 怎样证明结果可信?保存源位点、窗口版本向量、输入清单、三类 activation epoch、快照和风险动作,固定上下文后抽样回放对账。
  13. 话题合并为什么不能只按 event time 或归一后 offset 切换?event time 会迟到,topic 又会改变下游分区;所有权必须由归一前稳定日志的 barrier offset 唯一划分。
  14. 多窗口算分为什么不能只存一个 aggregate version?各聚合行独立推进,清单还要引用包含实际值和 source cut 的不可变版本,才能在 head 前进后重现。
  15. run 有 candidate count 为什么还可能少写?少一项再重复另一项时数量可以相同;既要先证明所有预期候选分区都在同一 source cut 提交 Top K,再冻结输入 manifest 精确集合,并让每项都落成候选或带原因的拒绝。
  16. 发布时为什么不推进 scope context version?发布只消费既定控制上下文;它条件锁定版本并 CAS 指针,只有规则、映射和风险控制事务负责推进 context。
  17. CDN purge 失败时紧急下榜怎么保证?命中后还要经过带 risk epoch 的边缘 denylist;边缘版本落后就绕过缓存,无法做边缘过滤时只能缩短 TTL 并在紧急状态绕过 CDN。
  18. 恢复代次切换后旧 worker 为什么不能靠更大的 aggregate_version 抢回 head?版本只在 generation 内比较,head CAS 同时校验 active generation 和当前版本,manifest 也引用完整三元组。

18.3 可选择的职责主线

如果项目真实存在,可把个人职责拆成互不冒领的方向:

  • 事件与流计算:负责契约、窗口、watermark、检查点和容量验证;
  • 算分与发布:参与规则版本、候选榜、快照原子发布和缓存恢复;
  • 风控与运营:协作主体去重、风险动作、控制台审计和紧急下榜。

只实现过 Redis 排行榜,就应明确讲到“当前方案”和可演进边界。没有真实事件规模、回放报告和监控证据,不能声称支撑过某个超大平台的生产热搜。

18.4 自检问题

  1. 热度和累计行为总量有什么区别?
  2. 为什么需要稳定 topic_id
  3. 新话题如何获得机会又不让垃圾词上榜?
  4. 去重主体为什么可能用近似算法?
  5. 同一个 actor 落到多个局部分区时怎样合并?
  6. watermark 和处理时间分别解决什么?
  7. 哪些迟到事件会修正在线结果?
  8. 规则版本如何与快照关联?
  9. 风险分与自然分为什么要分开?
  10. 快照发布为什么先写内容再切指针?
  11. 旧发布任务如何避免覆盖新快照?
  12. 某数据源中断时为什么不能照常套原公式?
  13. 候选榜丢失和正式榜丢失分别从哪里恢复?
  14. 哪些指标证明榜单没有被机器人主导?
  15. 什么时候批处理比流处理更合适?
  16. 话题合并时,barrier offset 和 effective watermark 分别解决什么?
  17. 一次候选算分读取多个窗口时,怎样证明没有混入不同水位的数据?
  18. 回滚到旧 rule_version 时,为什么仍要分配更高 activation epoch?
  19. 一个话题的风险动作如何让正在生成的旧 scope 快照失效?
  20. 窗口 head 前进后,旧 run 依靠什么重现当时的输入值?
  21. SEALED 怎样证明输入集合没有少处理一项,且封口后没有 worker 迟到写?
  22. 旧 normalizer 结果跨过合并 barrier 才到达时,怎样拒绝并重路由?
  23. 风险动作到期任务延迟时,为什么不能只按墙上时间自动解除?
  24. 一个候选分区整批未上报时,run 为什么不能仅凭已收到的 manifest 集合封口?
  25. aggregation generation 切换后,旧 worker 和旧 manifest 分别由什么 fence 阻止污染新代次?

总结

实时热搜榜的核心不是“给话题加分”,而是建立一条可解释、可恢复的实时决策链:

  1. 用可靠行为事件和稳定话题标识解决输入事实;
  2. 用多尺度窗口、主体去重、历史基线和时间衰减表达“此刻变热”;
  3. 用归一前稳定日志的分区 barrier、按 scope 的交接记录和映射 activation epoch 让话题合并只有一条事件负责链路;
  4. 用带 generation 的不可变窗口版本、候选分区 cut、输入清单和 run 级精确集合证明一次算分真正收齐、读取并处理了哪些输入;
  5. 用 scope 风控纪元、显式到期事务和边缘 denylist 控制刷量、缓存与内容安全;
  6. 用可封口候选 run、数据库权威指针和不可变快照分离持续计算与稳定发布;
  7. 用事件位点、窗口聚合、快照和隔离回放处理故障与争议。

Redis ZSET 可以是候选排序组件,但原始事件、窗口指标和已发布快照才让系统在重复、迟到、规则变化、热点和缓存故障后仍然说得清、恢复得了。