公司动态
热搜系统设计实战:基于Kafka与Redis的实时热榜实现
如果你做过后台开发大概率遇到过这样的需求产品经理拿着竞品截图说别人有热搜榜我们也要有。刚开始你会觉得这很简单不就是按关键词统计次数排个序吗但真正动手之后才发现这套系统远没有想象中那么轻松——数据从哪来、如何实时统计、连词和同义词怎么合并、旧话题怎么降下去、如何过滤广告和低俗内容、怎么防刷榜……任何一个环节没想清楚最后做出来的可能不是“热搜榜”而是一个谁都不愿意看的“Bot 灌水榜”。本文要讨论的就是“热搜话题系统”Trending Topic System的系统设计。我会从概念、架构、热度算法、代码实现、验证方式、常见问题和工程建议这些维度展开尽量用完整代码和配置带你跑通一个最小可用的热搜系统。读完这篇文章你不仅能应付系统设计面试里的相关题目更重要的是能在自己的项目中真正把热搜功能落地。1. 这篇文章真正要解决的问题搜索引擎和社交平台的热搜榜本质上是把“短时间内大量用户关注的话题”用一套可量化的规则挑选出来再按某种顺序展示给用户。但要设计好这套机制先要回答几个更底层的问题到底按什么指标衡量“热度”同一件事在不同平台上的关键词不一致时怎么聚合新话题刚冒头怎么识别旧话题长期霸榜怎么处理这些问题归结起来就是热搜系统设计中的四个核心难点实时性话题热度变化非常快统计和排序有延迟热搜榜就不“热”了。准确性原始数据是海量、稀疏、含噪声的文本流必须经过清洗、分词、聚类才能得到有意义的话题。公平性不能因为某个话题被少量 Bot 高频请求就自动冲上榜首要能识别和压制刷榜行为。平滑性榜单不能每秒都在剧烈抖动头条位置频繁变化会严重影响用户体验。从工程实现来看一个完整的热搜系统通常涉及数据接入、流式处理、在线存储、热榜计算、API 服务和前端展示等模块。它不是一个单机程序而是一套需要配合消息队列、实时计算引擎、Redis 等中间件完成的实时数据链路。这篇文章的实操部分会基于 Java Kafka Redis 搭建一个简化的热搜系统并用 Flink 风格的流处理思路说明实时聚合的实现原理。阅读本文你至少能收获三样东西搞清热搜系统整体架构和核心数据流掌握热度打分模型的常见设计和公式写法拥有一套可以本地运行的最小代码示例能直接改造成自己的原型。2. 热搜系统的核心概念与设计目标2.1 什么是热搜话题系统热搜话题系统是指从大规模用户行为数据搜索、点击、转发、评论、点赞、浏览中通过实时或准实时的方式识别出当前受关注度显著上升的话题并以榜单、推送、词云等形态呈现给用户的技术系统。它不只是“热门关键词排行榜”还可以是“热门事件”“热门问题”“热门商品”等不同粒度的热度对象。在系统设计层面热搜系统的输入通常是海量行为日志和内容文本输出是带有热度值、话题描述、时间戳等信息的榜单数据。中间发生的核心操作是采集、清洗、分词、聚合、打分、排序、缓存、下发。2.2 热搜系统的三类常见形态形态典型代表数据粒度更新频率关键词热搜搜索引擎热搜榜搜索词分钟级话题热搜社交媒体话题榜话题标签、事件聚类分钟级到小时级内容热搜短视频/资讯热榜单条内容分钟级不同形态的热搜系统技术难点并不完全相同。关键词热搜重点在于搜索日志的实时清洗和分词话题热搜重点在于把不同表述聚合到同一话题下内容热搜则更关注内容质量的筛选和个性化排序。本文的示例会以关键词热搜为主同时也说明如何扩展成话题热搜。2.3 热搜系统的设计目标从产品和技术两个维度看热搜系统要同时满足以下目标低延迟从事件发生到话题上榜目标通常在分钟级部分场景要求秒级。高可用榜单是核心流量入口不能因为热搜服务故障导致主站页面异常。可解释运营和用户都需要知道某个话题为什么上榜热度分要能拆解到各维度的贡献。可干预系统需要提供白名单、黑名单、人工置顶、人工撤销等运营手段。防刷防噪能识别爬虫、刷量行为和低质内容。成本可控海量日志的存储与计算要尽量低成本不能每个热搜词都做复杂的全量文本聚类。理解这些目标之后你会发现热搜系统的核心其实不只是“排序”而是“如何在噪音巨大的实时数据流中找到真正值得被看见的话题”。3. 热搜系统的整体架构与数据流3.1 架构总览典型的热搜系统分为五层数据源层 - 接入层 - 处理层 - 存储层 - 服务层数据源层用户搜索日志、点击行为日志、内容发布事件、业务数据库变更等。接入层通过埋点 SDK 或业务服务上报日志统一投递到消息队列。处理层实时计算引擎消费消息执行数据清洗、分词、频次统计、热度打分、聚合归类。存储层Redis 保存热榜结果和计数窗口ClickHouse/ES 保存明细数据和离线分析结果。服务层通过 HTTP API 或 SDK 向业务方提供热榜查询、话题详情、运营干预接口。3.2 核心数据流以最常见的搜索热搜为例一条完整的数据流如下用户在 app 或 Web 端输入关键词并搜索。前端或服务端通过埋点 SDK 上报搜索行为事件。日志进入 Kafkatopic 名为search_log。Flink 或 Spark Streaming 消费search_log执行清洗和分词。实时任务把“关键词 - 计数”的聚合结果写入 Redis。定时任务或流式任务基于 Redis 中的窗口计数计算热度分。热度分写入 Redis ZSet供 API 直接拉取排行。运营后台可以修改黑名单、白名单和置顶列表干预最终展示结果。从数据流向可以看出热搜系统本质上是一条“实时数仓”链路。离线部分可以基于 Hive 或 Spark 做历史话题分析为热度模型调参提供依据。3.3 为什么选择 Kafka Redis 作为主链路很多做小规模热搜系统的人会疑惑是不是一定要用到 Flink 和 Kafka 这么重的组件其实不是。选择哪套技术栈完全取决于数据规模和实时性要求如果每秒搜索日志只有几百条用 Java 的 ConcurrentHashMap 加定时线程池就能做热榜。如果每秒日志量达到几十万甚至上百万就必须靠 Kafka 削峰、Flink 分布式计算、Redis 高性能读取。本文示例选择了 Kafka Redis 的组合作为主链路是因为它在中小规模系统中足够简单同时又能把实时处理的思路讲清楚。实际生产环境再根据数据量补充 Flink 计算节点即可。4. 热度算法设计为什么不能只按次数排序如果只统计关键词出现次数你会看到满屏都是“疫情”“天气”“股票”这类常年高频词。它们不是“正在热”的话题而是“一直热”的话题。真正有用的热搜词是短时间窗口内热度上升最快的词。4.1 热度应该由哪些因素构成一个合理的热度分至少应该考虑以下因素因素说明获取方式搜索次数用户主动搜索某词的数量搜索日志内容产生量包含该词的帖子/视频数量内容发布事件互动量点赞、评论、转发等行为互动事件时间衰减距离当前时间越近权重越高计算得出用户权重高活跃用户的贡献更高用户画像在大多数场景下搜索次数和时间衰减是最核心的两个因素其余因素作为加权项来调节。4.2 时间衰减模型热度必须随时间衰减否则旧话题永远不会消失。常见的衰减方式有两种。一种是指数衰减。假设事件发生时刻为 (t_0)当前时刻为 (t)热度贡献值等于初始权重乘上 (e^{-\lambda(t - t_0)})。这种方式衰减速度快能快速反映热点变化适合对实时性要求高的场景。另一种是滑动窗口。只统计最近 N 分钟内的数据例如统计每个关键词在过去 10 分钟内的搜索次数窗口外的数据直接丢弃。这种方式实现简单但不适合对“突然爆发”和“缓慢上升”做精细区分。实际系统通常采用“滑动窗口 指数衰减”的混合方案用滑动窗口控制基础统计范围再用衰减系数对窗口内的较早数据降权。4.3 Hacker News 风格的评分公式这里给出一个可以落地的基础评分公式它参考了 Hacker News 和 Reddit 的排序思路但做了简化score (search_count 2 * content_count 3 * interaction_count) / (age_hours 2)^1.5其中search_count是时间窗口内的搜索次数content_count是时间窗口内产生的内容条数interaction_count是时间窗口内的互动次数age_hours是话题首次出现到现在的小时数。这个公式的特点是短时间内获得大量行为的话题会快速上升而随时间推移即使行为量不变分数也会逐步下降。你可以根据自己的业务调整分子中各项的权重和分母的指数。4.4 话题聚类如何把“新冠”和“新冠肺炎疫情”合并关键词热搜只需要分词和归一化即可但话题热搜还要求把同一事件的不同表述聚合到同一话题下。常见做法是文本归一化把大小写、全角半角、繁体简体统一。同义词扩展用词向量或编辑距离找出相似词。实体链接把“苹果发布会”“iPhone 16 发布”“苹果秋季发布会”都链接到同一实体或同一事件 ID。聚簇计算对高频词做 Embedding然后按向量相似度聚类再由聚类中心词作为话题名称。本文的代码示例只实现文本归一化和基础聚合更复杂的语义聚类需要引入 NLP 组件不在最小原型范围内展开。5. 环境准备与前置条件5.1 运行环境本文的示例代码基于 Java 17 Spring Boot 3 Kafka Redis使用 Docker Compose 启动中间件。如果你的机器没有 Docker也可以直接安装 Kafka 和 Redis但配置方式需要自行调整。5.2 需要安装的软件软件版本建议用途JDK17 或 21运行 Java 示例Maven3.8构建项目Docker最新稳定版启动 Kafka / RedisKafka3.x消息队列Redis7.x热榜存储和计数为了方便起见这里使用 Docker Compose 一次性启动 Kafka 和 Redis。你只需准备一个docker-compose.yml文件。version: 3 services: redis: image: redis:7-alpine container_name: trending-redis ports: - 6379:6379 kafka: image: bitnami/kafka:3.6 container_name: trending-kafka ports: - 9092:9092 environment: KAFKA_CFG_NODE_ID: 0 KAFKA_CFG_PROCESS_ROLES: controller,broker KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 0kafka:9093启动命令docker compose up -d5.3 项目依赖在pom.xml中加入 Spring Boot、Spring Kafka 和 Spring Data Redis 依赖dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.kafka/groupId artifactIdspring-kafka/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency dependency groupIdorg.apache.commons/groupId artifactIdcommons-text/artifactId version1.11.0/version /dependency如果你的项目是 Spring Boot 3.xcommons-text主要用于做文本归一化时的特殊字符处理也可以不要只用 JDK 自带的方法。6. 核心流程拆解从日志采集到热榜输出这一节先讲清楚热搜系统的实现步骤再在下一节给出可直接运行的完整代码。整个核心链路分为七个步骤第一步模拟搜索日志产生。真实场景中这一步由客户端和业务服务完成。日志内容包括用户 ID、搜索词、时间戳、来源渠道等字段。第二步日志投递到 Kafka。使用 Spring Kafka 的KafkaTemplate把日志消息发到search_logtopic。这里要注意消息 key 的设置建议用“会话 ID”保证同一用户的连续操作有序。第三步消费并清洗日志。消费者从 Kafka 拉取日志执行关键词提取。搜索词通常是用户输入的原样内容但需要做去空格、转小写、过滤标点等处理。第四步窗口计数聚合。将清洗后的关键词写入 Redis。最简单的方式是INCR自增但热搜系统要求的是时间窗口计数所以要用更精细的存储结构。本文示例会用“分钟级 Hash 定时合并”的方式实现窗口计数。第五步计算热度并写入 ZSet。定时任务扫描窗口计数计算每个关键词的热度分写入 Redis 的ZSet结构中以分数作为排序依据。第六步过滤与运营干预。从 ZSet 中取 Top N 时先过滤黑名单词再合并人工置顶词。第七步API 输出热榜。提供一个 HTTP 接口返回热榜带score和排名。这七个步骤中最容易出错的是第四步和第五步。窗口计数如果设计不好要么内存占用过高要么统计不准确。我推荐的方式是以“分钟”为最小粒度记录每个关键词每分钟的计数然后由定时任务聚合最近 N 分钟的数据。7. 完整示例Java 实现热搜系统最小原型本节以 Spring Boot 项目为例给出可以实际运行的热搜系统最小原型。代码文件尽量保持精简方便你直接复制到项目中修改。7.1 定义搜索日志消息体// 文件路径src/main/java/com/example/trending/model/SearchLog.java package com.example.trending.model; public class SearchLog { private String userId; private String keyword; private long timestamp; public SearchLog() {} public SearchLog(String userId, String keyword, long timestamp) { this.userId userId; this.keyword keyword; this.timestamp timestamp; } public String getUserId() { return userId; } public void setUserId(String userId) { this.userId userId; } public String getKeyword() { return keyword; } public void setKeyword(String keyword) { this.keyword keyword; } public long getTimestamp() { return timestamp; } public void setTimestamp(long timestamp) { this.timestamp timestamp; } }这里使用的 POJO 类是一个标准 JavaBeanJackson 可以直接把它序列化成 JSON 发送到 Kafka。7.2 日志生产者模拟用户搜索行为// 文件路径src/main/java/com/example/trending/producer/SearchLogProducer.java package com.example.trending.producer; import com.example.trending.model.SearchLog; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import java.util.Random; Component public class SearchLogProducer { private static final String TOPIC search_log; private final KafkaTemplateString, String kafkaTemplate; private final Random random new Random(); private static final String[] KEYWORDS { 系统设计, Redis, 消息队列, Java, 热搜系统, 大数据, Flink, Kafka, 微服务, Spring Boot }; public SearchLogProducer(KafkaTemplateString, String kafkaTemplate) { this.kafkaTemplate kafkaTemplate; } Scheduled(fixedRate 500) public void sendSearchLog() { String keyword KEYWORDS[random.nextInt(KEYWORDS.length)]; SearchLog log new SearchLog(user_ random.nextInt(1000), keyword, System.currentTimeMillis()); kafkaTemplate.send(TOPIC, log.getUserId(), toJson(log)); } private String toJson(SearchLog log) { return {\userId\:\ log.getUserId() \,\keyword\:\ log.getKeyword() \,\timestamp\: log.getTimestamp() }; } }这个生产者通过Scheduled注解每 500 毫秒向 Kafka 发送一条模拟搜索日志。实际生产环境中你不需要自己构造日志而是接入客户端埋点或服务端切面日志。7.3 消费者清洗日志并写入 Redis 中间计数// 文件路径src/main/java/com/example/trending/consumer/SearchLogConsumer.java package com.example.trending.consumer; import com.fasterxml.jackson.databind.ObjectMapper; import com.example.trending.model.SearchLog; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; import java.time.LocalDateTime; import java.time.format.DateTimeFormatter; Component public class SearchLogConsumer { private final StringRedisTemplate redisTemplate; private final ObjectMapper objectMapper new ObjectMapper(); private static final DateTimeFormatter MINUTE_FORMAT DateTimeFormatter.ofPattern(yyyyMMddHHmm); public SearchLogConsumer(StringRedisTemplate redisTemplate) { this.redisTemplate redisTemplate; } KafkaListener(topics search_log, groupId trending-group) public void onMessage(String message) throws Exception { SearchLog log objectMapper.readValue(message, SearchLog.class); String keyword normalize(log.getKeyword()); String minuteKey search:count: LocalDateTime.now().format(MINUTE_FORMAT); redisTemplate.opsForHash().increment(minuteKey, keyword, 1); } private String normalize(String keyword) { // 转小写、去除首尾空格和特殊字符 return keyword.trim().toLowerCase().replaceAll([\\p{Punct}], ); } }这个消费者使用 Hash 结构存储每个分钟窗口的关键词计数。search:count:202504101030这样的 key 表示 2025 年 4 月 10 日 10 点 30 分这一分钟field 是关键词value 是计数。这样做的好处是一个 minute key 只包含一分钟的数据聚合历史窗口时非常方便也方便定期清理过期 key。7.4 热度计算与榜单生成热度计算是热搜系统的核心逻辑。下面的TrendingScoreService会扫描最近 10 分钟的 Hash 计数计算每个关键词的热度分并写入 Redis ZSet。// 文件路径src/main/java/com/example/trending/service/TrendingScoreService.java package com.example.trending.service; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import java.time.LocalDateTime; import java.time.format.DateTimeFormatter; import java.util.HashMap; import java.util.Map; import java.util.Set; Service public class TrendingScoreService { private final StringRedisTemplate redisTemplate; private static final DateTimeFormatter MINUTE_FORMAT DateTimeFormatter.ofPattern(yyyyMMddHHmm); private static final int WINDOW_MINUTES 10; public TrendingScoreService(StringRedisTemplate redisTemplate) { this.redisTemplate redisTemplate; } Scheduled(fixedRate 30000) public void computeScores() { MapString, Long keywordCounts aggregateWindow(); String zsetKey trending:keywords; redisTemplate.delete(zsetKey); for (Map.EntryString, Long entry : keywordCounts.entrySet()) { double score computeScore(entry.getKey(), entry.getValue()); redisTemplate.opsForZSet().add(zsetKey, entry.getKey(), score); } } private MapString, Long aggregateWindow() { MapString, Long totals new HashMap(); for (int i 0; i WINDOW_MINUTES; i) { String minuteKey search:count: LocalDateTime.now().minusMinutes(i).format(MINUTE_FORMAT); SetObject fields redisTemplate.opsForHash().keys(minuteKey); if (fields null) { continue; } for (Object field : fields) { Object value redisTemplate.opsForHash().get(minuteKey, field); if (value ! null) { long count Long.parseLong(value.toString()); totals.merge(field.toString(), count, Long::sum); } } } return totals; } private double computeScore(String keyword, long searchCount) { // 简化热度公式搜索次数 / (已存在小时数 2)^1.5 double ageHours getAgeHours(keyword); return searchCount / Math.pow(ageHours 2, 1.5); } private double getAgeHours(String keyword) { // 实际项目中记录关键词首次出现时间这里用固定值演示 return 1.0; } }这里的热度计算做了一次简化使用固定的ageHours值演示公式。实际项目中你需要把“关键词首次出现时间”存入 Redis再从当前时间减去首次出现时间得到ageHours。7.5 查询接口返回热门榜单// 文件路径src/main/java/com/example/trending/controller/TrendingController.java package com.example.trending.controller; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.data.redis.core.ZSetOperations; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestParam; import org.springframework.web.bind.annotation.RestController; import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.Set; RestController public class TrendingController { private final StringRedisTemplate redisTemplate; public TrendingController(StringRedisTemplate redisTemplate) { this.redisTemplate redisTemplate; } GetMapping(/api/trending) public ListMapString, Object getTrending( RequestParam(defaultValue 10) int topN) { String zsetKey trending:keywords; SetZSetOperations.TypedTupleString tuples redisTemplate.opsForZSet().reverseRangeWithScores(zsetKey, 0, topN - 1); ListMapString, Object result new ArrayList(); int rank 1; if (tuples null) { return result; } for (ZSetOperations.TypedTupleString tuple : tuples) { result.add(Map.of( rank, rank, keyword, tuple.getValue(), score, tuple.getScore() null ? 0 : Math.round(tuple.getScore() * 100.0) / 100.0 )); } return result; } }这个接口使用reverseRangeWithScores从 ZSet 中取出分数最高的 N 个关键词并带上排名和分数返回 JSON 给前端。7.6 启动类和配置最后需要配置application.yml中关于 Kafka、Redis 和定时任务的参数。# 文件路径src/main/resources/application.yml server: port: 8080 spring: kafka: bootstrap-servers: localhost:9092 consumer: group-id: trending-group auto-offset-reset: earliest key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer data: redis: host: localhost port: 6379 # 启用定时任务 spring.task.scheduling.enabled: true// 文件路径src/main/java/com/example/trending/TrendingApplication.java package com.example.trending; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.scheduling.annotation.EnableScheduling; SpringBootApplication EnableScheduling public class TrendingApplication { public static void main(String[] args) { SpringApplication.run(TrendingApplication.class, args); } }8. 运行结果与效果验证8.1 启动步骤按照下面的顺序启动整个系统# 1. 启动 Kafka 和 Redis docker compose up -d # 2. 启动 Spring Boot 应用 mvn spring-boot:run启动日志中如果没有报错说明 Kafka 和 Redis 连接成功。生产者会每 500 毫秒生成一条日志消费者持续消费并写入 Redis。8.2 验证 Kafka 数据是否正常可以先进入 Kafka 容器查看 topic 是否创建成功docker exec -it trending-kafka kafka-topics.sh --bootstrap-server localhost:9092 --list如果看到search_log说明 topic 创建成功。如果需要手动创建 topic可以执行docker exec -it trending-kafka kafka-topics.sh --bootstrap-server localhost:9092 --create --topic search_log --partitions 1 --replication-factor 18.3 验证 Redis 计数和榜单先查看分钟计数docker exec -it trending-redis redis-cli KEYS search:count:*应该能看到最近几分钟的 Hash key。再查看某个 key 的内容HGETALL search:count:202504101030等 30 秒让定时任务运行后查询热榜ZREVRANGE trending:keywords 0 9 WITHSCORES也可以直接通过 HTTP 接口获取榜单curl http://localhost:8080/api/trending?topN10预期返回类似下面的 JSON[ { rank: 1, keyword: redis, score: 12.34 }, { rank: 2, keyword: java, score: 10.21 }, { rank: 3, keyword: 系统设计, score: 8.02 } ]如果返回的榜单与预期不符优先检查三个方面Kafka 消费者是否消费到了消息、Redis 中是否有分钟计数、定时任务是否执行成功。排查顺序建议是 Redis - Kafka - 消费者代码。9. 热搜系统的常见问题与排查思路问题现象可能原因排查方式解决方案榜单为空Redis ZSet 为空查看 Redis 是否有trending:keywordskey检查定时任务是否执行检查消费者是否写入计数榜单不更新定时任务未执行查看日志中是否有computeScores执行记录确认EnableScheduling已开启检查 Cron/固定频率配置关键词大小写不合并清洗逻辑未归一化查看 Redis Hash 中的原始 field在消费者清洗阶段统一转小写同一个话题出现多个相似词没有做同义词或聚类观察榜单中是否有近义词引入同义词表或 Embedding 聚类统计延迟分钟级窗口聚合逻辑运行较慢查看定时任务执行耗时优化 Redis 批量获取增加聚合频率改用 Flink 实时聚合旧热搜词长期霸榜时间衰减公式不合理检查热度计算公式中ageHours是否动态增长用首次出现时间计算年龄调大分母指数刷量词冲上榜首没有防刷策略查看异常用户请求频率增加用户权重、频控和 IP 维度的过滤逻辑9.1 关于分区和乱序的说明Kafka topic 的 partition 数量会影响消息处理的并发度和顺序性。在热搜系统中日志的全局顺序并不重要因此可以适当增加 partition 数量提升消费者并发度。但如果同一个用户短时间内多次搜索同一个词且你希望精确去重就需要按用户 ID 作为 key 发送消息保证同一用户的消息进入同一 partition。9.2 关于 Redis 内存膨胀问题分钟级别的 Hash 计数如果长期不清理Redis 内存会被撑爆。建议通过 Redis 的EXPIRE命令给每个 minute key 设置过期时间例如设置为 1 小时确保超过聚合窗口的就自动删除。生产环境还可以用定时任务统一清理或者直接使用 Redis Cluster 分片扩容。10. 热搜系统的工程化最佳实践10.1 使用 Flink 替代定时任务做真正实时的聚合本文示例用Scheduled定时任务模拟了流式聚合但它的实时性只能到秒级或分钟级且窗口聚合逻辑不灵活。生产环境更推荐使用 Flink SQL 或 DataStream API 直接做窗口聚合把结果写入 Redis。Flink 的优点是可以精确控制滑动窗口和滚动窗口并且在窗口内执行去重和防抖逻辑。下面是一个 Flink SQL 实现窗口计数的示意CREATE TABLE search_log ( user_id STRING, keyword STRING, ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL 5 SECOND ) WITH ( connector kafka, topic search_log, properties.bootstrap.servers localhost:9092, format json ); CREATE TABLE trending_result ( keyword STRING, cnt BIGINT, window_start TIMESTAMP(3), PRIMARY KEY (keyword, window_start) NOT ENFORCED ) WITH ( connector upsert-kafka, topic trending_result, properties.bootstrap.servers localhost:9092, format json ); INSERT INTO trending_result SELECT keyword, COUNT(*) AS cnt, TUMBLE_START(ts, INTERVAL 5 MINUTE) AS window_start FROM search_log GROUP BY keyword, TUMBLE(ts, INTERVAL 5 MINUTE);这段 SQL 的作用是每 5 分钟统计一次每个关键词的搜索次数并把结果写入 Kafka 的trending_resulttopic。下游服务消费这个 topic 再写入 Redis 榜单即可。需要说明的是这里保留窗口开始时间是为了方便做多级聚合实际写 Redis 时可以只保留最新窗口结果。10.2 榜单缓存与降级策略热搜榜是高频读取、低频更新的数据必须做缓存。最简单的做法是使用 Redis 本身作为存储不需要再额外增加本地缓存。API 服务查询时可以先查本地 Caffeine 缓存命中不到再查 Redis或者直接直接用 Redis 的 ZSet 查询设置较短的过期时间。对于极端情况下的 Redis 故障API 服务需要返回兜底数据——例如最近一次成功生成的榜单快照而不是直接 5xx。10.3 防刷策略必须前置防刷不能只靠热度算法去“事后降权”必须在数据进入时就过滤。常用策略包括用户维度频控同一用户在短时间内超过阈值直接丢弃或降权。设备/IP 维度频控识别机刷流量。文本质量过滤过滤纯数字、无意义重复词、广告词。黑名单动态更新运营可以在后台实时添加黑名单词消费者侧同步过滤。这些策略可以放在消费者入口处也可以在 Flink 作业里作为 Filter 算子实现。原则是越早过滤越省计算资源。10.4 运营干预接口设计纯算法生成的榜单不一定符合产品预期所以运营干预是标配能力。建议提供以下能力关键词置顶指定词固定排在第一位或前几位。关键词加白即使热度不够也能上榜。关键词屏蔽词语永远不会出现在榜单中。热度加权或降权人工调整某个词的分数倍数。从架构上讲运营干预配置一般存在 MySQL 或 Apollo 配置中心中榜单查询时实时读取并应用规则。为了性能可以定期把配置加载到 Redis。10.5 监控与指标热搜系统至少需要监控以下指标日志消费延迟Kafka 消费 lag 是否持续增加。Redis ZSet 大小是否为预期规模是否存在异常膨胀。定时任务执行时间是否超过窗口间隔。API 查询延迟P99 是否在可接受范围内。更新频率榜单每次更新的间隔是否稳定。生产环境可以接入 Prometheus Grafana把这些指标做成看板。热搜系统一旦出现故障用户是不会在页面上看到报错的只会看到“榜不动了”或者“榜很奇怪”。监控的意义在于提前发现异常而不是等用户反馈。11. 总结与后续学习方向热搜系统表面上是一个“实时排行榜”功能但把它拆开看涉及的子问题涵盖了数据采集、文本清洗、流式计算、窗口聚合、热度建模、缓存设计、防刷策略、运营干预等多个领域。这也是为什么它经常出现在系统设计面试题目中——它非常考验设计者对实时数据链路的整体理解和工程权衡能力。本文围绕热搜系统的主要技术难点展开从架构、数据流、热度公式到 Java 原型代码基本覆盖了从零搭建一个简化版热搜系统的完整过程。建议你先在本地把示例代码跑通再尝试扩展几个方向把定时聚合改成 Flink 实时聚合体验窗口计算的差异引入同义词表和实体链接把关键词热搜升级成话题热搜增加用户权重和防刷逻辑看热度排序会发生什么变化把榜单数据改造成基于 WebSocket 的实时推送模拟 App 端下拉刷新效果。在动手实践时有一个容易忽略的设计点要提醒你热度算法的参数一定不是上线前定死的而是要根据线上数据持续调优。不同业务场景下搜索、内容、互动三类行为的权重差异很大。电商平台可能更看重加购和成交内容社区可能更看重转发和评论搜索引擎可能更看重搜索频次和搜索用户数。先做一套基础公式上线再通过 A/B 实验调参往往比一开始追求复杂模型更实用。另外如果你是在学习系统设计建议把这篇内容与“经典系统设计面试题”一起对比复习设计 Twitter Trending、设计微博热搜、设计抖音热榜这些题目虽然是不同产品形态但底层逻辑高度一致。理解了一个其余都能触类旁通。热搜系统的核心从来不只是统计和排序而是定义“什么值得被看见”。想清楚产品意图再选择合适的技术方案才是这套系统真正有价值的地方。