
做实时新闻热搜这个项目我最大的感受是热榜不是“排”出来的是“算”出来的。你手机端那个每秒钟都在跳的“实时热榜”背后其实是一条从点击事件采集、消息队列缓冲、Flink 流式计算到榜单存储和可视化展示的完整链路。这篇文章就把我这个实战项目完整拆一遍从业务理解到架构选型从窗口计算到热度衰减算法的落地从集群部署到上线后排查过的几个经典故障尽可能把可以直接抄作业的部分讲透。这个项目的受众很明确正在做大数据方向毕业设计的学生、准备大数据开发岗位面试的人以及已经接触过离线的 MapReduce、Spark 批处理想往 Flink 实时计算方向转的开发者。如果你是其中任何一类这篇内容基本能帮你少走一个月的弯路。我会把每个关键选择背后的“为什么”讲清楚而不只是贴一段代码让你复制。1. 项目定位实时新闻热搜到底要解决什么问题1.1 先和“产品”把榜单口径对齐接到“新闻热搜”这类项目最大的坑不是技术而是需求边界模糊。你要是上来就闷头写 Flink 任务大概率做完发现产品和用户要的根本不是一回事。我在动手之前把实时热搜拆成了四个明确指标第一个指标是统计周期内的点击量增量。比如过去 5 分钟内某条新闻被点击了多少次这决定了一条新闻“当前”火不火。第二个指标是滑动窗口内的综合热度比如过去 2 小时内反复点击、评论、转发的新闻不能因为某一分钟波动就被挤下去。第三个指标是分类热榜娱乐、体育、科技要分开看否则科技类新闻永远被娱乐新闻压在下面。第四个指标是突发新闻识别也就是某条新闻的点击增速短时间内出现陡增这种“爆发性”才是热搜最有价值的部分用户看到的是“新”而不是“旧”。这四层需求对应到技术实现上前三层靠窗口聚合就能完成第四层则需要结合增速计算和阈值判断。没有这些前置拆解你就不知道 Flink 作业里该开几个窗口也不知道 Redis 里到底要存什么结构的数据。1.2 为什么选 Flink 而不是 Spark Streaming 或 Storm在这类项目里常见的技术选型对比是 Flink、Spark Streaming 和 Storm。Storm 的延迟确实低但是它的实时计算模型太原始没有内建的状态管理要自己处理状态备份和故障恢复开发成本和运维成本都很高现在新项目已经很少用它了。Spark Streaming 其实是微批处理默认每两秒到几秒提交一批虽然吞吐高但批与批之间存在调度延迟而且它的精确一次语义实现相对绕。在新闻热搜场景里热点事件可能在一两分钟内爆发微批的模式会明显感觉到榜单刷新有“顿挫感”。Flink 是真正的流式计算引擎事件一进 KafkaFlink 就能立刻处理毫秒级延迟。另一个决定性优势是 Flink 有内建的窗口机制和状态管理。窗口可以做滚动、滑动、会话窗口状态可以自动做 checkpoint配合 Kafka 能实现精确一次语义。对于新闻榜单这种需要长时间累计状态、且要求秒级刷新的场景Flink 确实是当前最顺手的选择。1.3 项目的输入、输出和评价标准明确了指标还要框定项目的输入输出边界。这个项目的输入是模拟的新闻点击事件流我把它定义成一条 JSON 格式的埋点日志包含新闻 ID、标题、分类、用户行为类型点击/评论/分享、事件时间戳等字段。输出则是三层实时热榜前 50 条、分类热榜前 20 条、突发新闻列表。分别对应 Redis 中的 ZSET、Hash 和一个独立列表。评价标准上比较重要的是时效性。从事件发生到榜单更新端到端延迟目标设定在 10 秒以内实际上 Flink 处理这部分只占几十毫秒延迟大头通常在消息缓冲和前端轮询间隔上。另一个标准是准确性窗口边界重叠时不能重复计数同一个用户的重复点击要有去重机制。这两点在后文的窗口设计和状态编程里都会体现。2. 整体架构设计与技术选型解析2.1 用“四层架构”来规划整个系统大数据项目习惯上会讲四层架构数据采集层、数据计算层、数据存储层、数据应用层。这个分类在这个项目里同样适用。采集层负责把新闻客户端的埋点日志实时收集起来这里我直接用 Java 程序模拟埋点发送到 Kafka计算层就是 Flink 集群消费 Kafka 数据完成窗口聚合、热度计算、关联维表存储层根据读写特征不同分别使用 Redis、MySQL可选的还有 Elasticsearch应用层则是后端查询接口加前端 ECharts 可视化大屏。做架构设计时最容易犯的错是总想着用一套存储搞定所有需求。热度榜需要毫秒级读取用 Redis历史榜单需要按时间回溯统计放 MySQL 更稳全文检索和词频分析可以挂 ES。各司其职比强行造一个大而全的存储要务实得多。2.2 数据链路全景和各节点职责整条链路是模拟新闻服务端埋点程序发消息到 Kafka 的news_click_log主题Flink 作业以消费者组news-hot-group订阅这个主题经过事件时间窗口聚合后把“每 1 分钟各新闻的增量指标”写入 Redis 集群同时异步写一份到 MySQL保存历史明细另一个 Flink 作业或者同一个作业的后续算子负责从 Redis 读取累计热度计算榜单前 N 名推给后端 API前端通过 WebSocket 或者短轮询拿到数据用 ECharts 渲染实时曲线和词云。这里需要特别说明一个设计取舍。如果福东直接把 2 小时滑动窗口的热度结果写 Redis数据量还好但如果把全量新闻状态都存在 Flink 内部做热度衰减任务的状态后端会变得非常大而且一旦并行度调整重新分配状态的代价很高。所以我在实际项目里采用了一种配合方案Flink 只负责算每 5 分钟的增量指标把“衰减模型”放到 Redis 端和应用层去做增量累加。Flink 状态小、重启快Redis 里的 score 也是累加的两边职责非常清晰。2.3 关键组件选型对照表选型这件事很多教程会直接告诉你用哪个但我更想把对比逻辑写出来让你面对其他项目时也能自己判断。组件维度上消息队列我在 Kafka 和 Pulsar 之间选。Kafka 生态最成熟Flink 连接器支持最稳定吞吐和消息回溯能力足够Pulsar 的架构更现代化但对这个项目来说优势体现不出来。Flink 与 Kafka 的组合在社区里几乎是标准答案遇到问题搜解决方案也最快。存储方面实时榜单用 Redis 的 ZSET因为 ZINCRBY 命令天然支持对一个 key 的 score 累加获取 Top N 用 ZREVRANGE 就是一次 Redis 请求性能极高。Redis 内部需要定期清理过期 key所以我给每条新闻的热度 key 设置了过期时间配合滑动窗口继续重建。MySQL 用来落历史数据表结构比较简单核心字段是新闻 ID、标题、分类、统计时间、点击量、热度分。如果后续想实现“按天回看历史热搜榜”MySQL 是最方便的。另外这个项目还挂了一个可选的 Elasticsearch用来做新闻标题的词频统计和热搜词云。Flink 可以通过 Elasticsearch connector 把分词后的词频结果写入 ES前端用 agg 聚合查询生成词云。不是必须的模块但加上它整个项目在答辩或面试时会更完整。2.4 实时项目与离线批处理项目的差异补充我在做这个项目之前先做过基于 Spark 的离线网约车数据分析项目同样是处理“数据进、结果出”两者的思维差异还是很大的。离线项目是跑完一批算一批数据都在 HDFS 或者 Hive 表里算错重跑就行时间上不敏感。实时项目里的数据像流水一样不断进来窗口一关闭数据就过去了你想重算某一段数据常常发现已经来不及。因此在实时项目的开发规范里我养成了几个习惯所有输入数据必须带事件时间字段计算逻辑必须保证幂等下线重跑前要确认 Kafka 的 offset 清零否则重复数据会把榜单污染。3. 核心实时计算逻辑与关键代码实现3.1 事件协议与数据格式定义所有实时计算的起点都是数据结构设计。我定义的新闻点击事件长这样{ news_id: n10001, news_title: 某地发布高温橙色预警, category: weather, action: click, user_id: u98312, event_time: 2025-06-01 12:30:45, source: android }第一个要点event_time必须由业务端写入代表用户点击发生的真实时间。不能依赖 Flink 的PROCESS_TIME因为 Kafka 消费有延迟如果用处理时间计算窗口窗口边界会和真实时间错位榜单高峰时段会出现明显偏移。第二个要点action字段区分点击、评论、分享为后续多维度热度加权做准备。第三个要点user_id要参与去重逻辑同一个用户在一分钟窗口内只记一次点击否则一台设备刷几次就能把新闻刷上热搜。3.2 没有埋点数据怎么办自定义 DataSource刚搭环境时往往没有真实埋点数据源这时不要卡住先用 Flink 自定义 DataSource 模拟事件流。这也是很多 Flink 教程中“自定义 DataSource”的实际用途让你在没有 Kafka、没有外部数据的情况下也能跑通逻辑。自定义 Source 的核心是继承RichParallelSourceFunction在run()方法里循环构造格式正确的事件对象发送出去用Thread.sleep()控制发送速率。比如设置每 50 毫秒发一条事件、单条新闻随机权重就能模拟出某条新闻突然点击量陡增的“热点爆发”场景。这个 Mock Source 还能加上随机分类、随机新闻 ID用来测试窗口聚合逻辑是否正确。我自己写的时候加了可配置速率和过滤条件比如可以指定某条新闻发送 100 条/秒这样验证突发新闻识别逻辑时不需要等真实流量。整个 Source 只在测试环境跑上线换成 Kafka Source 即可代码主体的窗口计算逻辑一行都不用改。3.3 核心窗口计算滚动窗口产出增量指标窗口设计是整个项目的主心骨。项目里我用了两套窗口。第一套是 1 分钟滚动窗口产出“每分钟增量指标”它的主要作用是给热度衰减模型提供输入。Flink SQL 里创建 Kafka 源表后聚合逻辑这样写INSERT INTO mysql_daily_stat SELECT news_id, news_title, category, window_start, window_end, COUNT(DISTINCT user_id) AS click_cnt, COUNT(IF(action comment, 1, NULL)) AS comment_cnt, COUNT(IF(action share, 1, NULL)) AS share_cnt FROM TABLE( TUMBLE(TABLE news_click_log, DESCRIPTOR(event_time), INTERVAL 1 MINUTE) ) GROUP BY news_id, news_title, category, window_start, window_end这里有个很关键的操作COUNT(DISTINCT user_id)在 Flink SQL 里是精确去重但状态量随窗口内用户数线性增长。在高并发场景下我会改用一个 Sketches 近似基数计算来降低状态大小因为新闻榜单差几十个点击量对最终排名影响不大性价比很高。对于本项目的模拟数据量精确去重完全够用不用过度设计。第二套窗口是 1 小时滑动窗口每 5 分钟输出一次当前热度。这个窗口直接对应前端页面上看到的“最近一小时热榜”。滑动窗口的特点是窗口区间重复同一事件会被多个窗口包含Flink 会把重复计算处理好但这也意味着同样的数据会被重复输出多次下游写 Redis 时要用 UPSERT 语义防止重复累加。3.4 热度计算模型不能只看点击量新闻热搜不能简单按点击量排名否则热门事件会永远霸榜新事件根本挤不进来。热度分设计我用了一个加权衰减模型hot_score(t) (click_cnt * 0.5 comment_cnt * 1.5 share_cnt * 2.0) * decay_rate decay_rate exp(-lambda * (current_time - last_update_time))这个公式的意思是评论和分享的权重高于普通点击同时热度随时间指数衰减。每过 5 分钟应用层从 Redis 读取上一条新闻的旧热度分套这个公式更新一次。这样即使一条旧新闻持续收到新流量它的热度增长也会衰减突发新闻才有机会快速冲到前排。这个功能如果放在 Flink 里直接用状态实现也不难。用一个 KeyedProcessFunction注册定时器每 5 分钟触发一次读取当前值和上次更新把旧值的衰减部分加新值再更新状态。但落地时我发现收益不明确Redis 那边本来就是权威存储Flink 再存一份全量状态会增加先状态大小。最后我选择了把增量计算放在 Flink把衰减累加放在 Redis 的方式虽然多一次 Redis 读写但整体简单可靠。3.5 自定义 DataSinkRedis 写入的加速技巧既然要写 Redis 热榜最直接用的是 Jedis 客户端。要注意的是不要每条数据都建立一个新连接而是用连接池并且大批量写入时使用 Pipeline。比如一分钟内有一万条增量记录如果逐条写就是一万次网络往返用 Pipeline 可以压缩到一次批量提交耗时能降到原来的十分之一。public static class RedisSink extends RichSinkFunctionTuple3String, Double, Long { private JedisPool jedisPool; Override public void open(Configuration parameters) { jedisPool new JedisPool(new JedisPoolConfig(), 127.0.0.1, 6379); } Override public void invoke(Tuple3String, Double, Long value, Context context) { try (Jedis jedis jedisPool.getResource()) { Pipeline pipeline jedis.pipelined(); pipeline.zincrby(hot_rank_latest, value.f1, value.f0); pipeline.hset(news_meta, value.f0, String.valueOf(value.f2)); pipeline.sync(); } } }写 Redis 时有几个容易踩的坑。首先是 key 设计不同时间窗口的榜单要用不同 key否则每分钟的任务会把上一分钟的数据覆盖。我的方案是把 key 设计成hot_rank_yyyyMMdd_HHmm这类带时间戳的形式方便按时间回看。然后是序列化value 里的 news_id 建议用原始字符串不要用 JSON 包一层越简单的结构在 Redis 里越好处理。最后是连接池配置maxTotal可以根据 Flink 写入并行度来设置一般设为并行度乘以 2 到 3 比较稳妥。4. 从零搭建环境Flink 与 Kafka 的部署实操4.1 版本选择与安装顺序这套环境的版本搭配我推荐 JDK 1.8 或 11、Kafka 2.13-2.8.0、Flink 1.14 或 1.16。Flink 1.14 的 Streaming SQL 已经比较成熟社区资料丰富排查问题方便。如果一定要用更新的版本建议对应查看 Flink 与 Kafka 连接器的兼容性表避免出现序列化版本不一致的报错。安装顺序也有讲究因为 Kafka 依赖 Zookeeper 做协调服务。推荐顺序是 JDK、Zookeeper、Kafka、Flink。每一步安装完先验证再继续不要一次装完一堆再排查。所有软件都装在同一台 CentOS 7 虚拟机或者用三台机器组成一个小集群都是可行的。三台机器时Zookeeper 配置成奇数节点如 3 个Kafka 的 broker 分别指向三个节点Flink 配一个 master 和两个 taskmanager这样能演示容错和并行度。4.2 Kafka 环境部署与主题创建Kafka 配置里重点要改的是server.properties的zookeeper.connect和log.dirs路径。启动顺序不能乱先启动 Zookeeper确认 2181 端口监听正常再启动 Kafka。启动完用生产者、消费者脚本各发一条消息验证一下一定确保基础通信正常再去做 Flink 对接。创建主题的命令bin/kafka-topics.sh --create \ --bootstrap-server node01:9092 \ --replication-factor 2 \ --partitions 3 \ --topic news_click_log主题分区数是这里最关键的参数。Flink 读取 Kafka 的并行度上限等于分区数一个分区对应一个并行子任务。如果你把并行度设为 10而分区只有 3那 7 个并行就空转了。因此建议主题分区数按生产预估吞吐量除以单分区吞吐来定本项目我设置为 3对应三点式模拟输入的并发足够支撑每天几千万级点击量。4.3 Flink 集群的 Standalone 模式部署如果是单机演示Standalone 模式最直接。修改conf/flink-conf.yaml的几个核心参数taskmanager.memory.process.size设大一点我设置成 2048mparallelism.default设为 2state.backend设为 rocksdbstate.checkpoints.dir指向本地或 HDFS 路径。如果是本机试验没有 HDFS也可以指到一个本地目录演示 checkpoint 已经足够。启动集群命令bin/start-cluster.sh启动后用 JPS 查看有没有出现 StandaloneSessionClusterEntrypoint 和 TaskManagerRunner 进程然后在浏览器打开http://node01:8081。如果能看到 Flink Web UI说明集群状态正常。这个 UI 对后面的故障排查非常重要反压、延迟、checkpoint 失败都能在上面直接看到。4.4 提交 Flink 作业并观察运行状态如果代码是用 Flink SQL 写的可以直接在 sql-client 里执行建表和插入语句。如果是 DataStream API需要打包成 jar 提交flink run -m node01:8081 \ -d \ -c com.example.NewsHotAnalyzer \ ./news-hot-analysis.jar提交后重点观察三个指标Kafka 的消费 Lag 是否持续增大、Task 的延时是否正常、Checkpoint 是否连续成功。我一般让作业先跑几分钟如果 Checkpoint 失败两三次就会立即排查绝不带病上线。这里有个经验如果 Checkpoint 总是超时多半是作业里用到了需要长时间持锁的操作比如事务型 Sink或者某个算子处理耗时过长优先级最高的排查方向一定是对齐阶段卡在哪个 Task 上。5. 上线后常见的坑与排查实录5.1 报错“JDBC 连接器异常”搜这个关键词的人非常多我也踩过。首要原因是缺少数据库驱动 JARFlink 官方 JDBC 连接器依赖一个 SPI 动态加载驱动的机制驱动 JAR 必须放在 Flink 的 lib 目录或者作业的 fat jar 里。报错往往是ClassNotFoundException或Unable to obtain JDBC Connection。第二个原因是 JDBC 连接参数问题。MySQL 8.x 连接串里serverTimezone必须设置否则连接建立时会报时区错误。Flink 的 JDBC Sink 默认使用批量提交参数里配置了sink.buffer-flush.interval以后如果 SQL 执行失败Flink 会重试重试次数耗尽后整个作业会 fail。排查这类问题时优先看 TaskManager 日志里的完整堆栈别只盯着 Web UI 上的报错摘要。5.2 想把结果写 Hive 表但数据一直不入表项目后期我尝试把历史榜单直接同步到 Hive 表做离线分析这也是网上提问很多的问题Flink sink 到 Hive 表数据不落地。核心原因有三个。第一是 Hive 表的分区路径没有自动创建或没写权限。Flink 写入 Hive 时如果目标分区在 HDFS 上已经存在会采用覆盖或追加策略如果不存在则有另建分区的开销在部分版本上会直接静默失败。排查时先去 HDFS 页面上看对应分区目录是否存在。第二是 Hive 的 ACID 事务表Flink 的写入器要拿到排他锁如果有人在同张表上跑了其他任务锁冲突就会让写不进去。第三是 Flink 的 Hive Sink 默认是批量提交数据会先落在临时目录只有触发 checkpoint 或设置sink.partition-commit.policy.kind才会 move 到正式目录。所以要确认 checkpoint 是开启的否则数据一直留在temp目录里看不到。5.3 热点爆发时消费 Lag 上涨和背压处理新闻热搜系统流量特征非常明显深夜平稳早晚高峰和突发新闻阶段流量会瞬间涨上几倍。这时 Flink 作业经常出现背压反压也就是下游处理不过来上游 Kafka 消费 Lag 持续上涨。处理背压先不能盲目加并行度。第一步是定位卡点。打开 Flink Web UI 的 BackPressure 页签查看哪些 Task 处于 HIGH 状态。通常问题集中在 Sink 节点因为我这个项目最常写 Redis如果 Redis 的 QPS 达到瓶颈整个链路就会阻塞。第二步是优化瓶颈算子。写 Redis 时检查连接池够不够大Pipeline 有没有真正生效。第三步才是考虑增加并行度。而且要注意增加并行度时源端 Kafka 的分区数必须也足够否则扩容毫无意义。5.4 数据质量检查不能省实时链路里最怕的就是脏数据。我专门写了一个轻量级的数据质量检查模块挂在 Flink 作业的分支上。核心检查项包括事件时间是否合法有没有超过当前时间的两小时时间戳乱跳的数据直接丢掉news_id 是否为空用户行为类型是不是合法枚举值。这类检查不是过滤而是双路输出一条流进正式分析链路另一条把校验失败的数据写到 Kafka 的重试主题或专门的坏数据主题里。这样才能在榜单出问题时快速定位是算法问题还是数据问题而不是对着结果空猜。项目上线第一周这个模块就帮我抓到了模拟数据程序中时间戳没有加时区、导致事件时间整整快 8 小时的 Bug。5.5 用火焰图定位 CPU 瓶颈如果作业能跑但吞吐上不去、CPU 使用率很高火焰图是非常好用的分析工具。可以通过 Java Flight Recorder 或者 async-profiler 对 TaskManager 进程采样生成火焰图。用命令直接抓./profiler.sh -d 60 -f /tmp/flink_flamegraph.html TaskManager-PID火焰图观测的重点是“平顶”和“宽底”。如果顶点是某个 JSON 解析工具说明数据反序列化开销太大考虑换更快的序列化方案如果某块区域很宽说明某个方法反复调用比如我在早期版本里每条事件都去构建一个 SimpleDateFormat 对象火焰图打开后这个方法的调用栈宽得离谱改成一个静态共享的 DateTimeFormatter 之后吞吐直接提升了三成。火焰图要结合 GC 日志一起看如果 GC 线程占用很高优先调内存参数而不是揪代码逻辑。6. 一点收尾的个人体会这个项目做完后的某一天我把 Flink 作业停了单独用 Spark 离线任务重新刷了一遍同样的历史数据对比实时榜单和离线榜单的差异。虽然业务上两者本来就是不同口径但这种对比让我真正理解了实时计算的定位它不适合做复杂的全量重算所有计算都必须在数据流经的几毫秒内完成因此设计时要想清楚什么是“提前算好”的什么是“真正增量”的。如果你也准备做类似项目我强烈建议不要一开始就去追求炫酷的技术栈而是把一个真实业务问题跑通。模拟流量可以先小一点把窗口、热度函数、存储、展示的每一环都弄清楚再逐步加复杂度和数据量。这套项目的底层逻辑扩展到电商实时热点、短视频热榜、舆情分析等场景基本是同一个套路把数据源换掉把指标定义改一改整个架构就能复用。踩过的坑已经写在上面剩下真正要踩的才是属于你自己的经验。