译者按
核心结论一句话:Hotstar 把表情服务从第三方迁到自建,靠“客户端异步缓冲 + Kafka + Spark 微批聚合 + PubSub 实时推送”这条流水线,用异步和批量换吞吐,扛住了单场赛事 50 亿次提交。
落地到国内中小厂得打个问号:这套架构暗含你已经有成熟的 Kafka 数据平台和 Spark 集群,K3s/K8s 下自己裸跑这两样运维成本直接劝退,常规场景用云上托管 Kafka 加 Redis 聚合完全够用,别照搬。
踩坑提醒一句:生产端“500ms 或 2 万条触发批量写入”的缓冲设计,突发流量下本地缓冲积压容易引发内存水位上涨,记得给进程加内存监控和积压告警;另外在 K8s 里跑 2 秒粒度的 Spark 微批,务必提前压测并预留 CPU 余量,否则 HPA 扩缩容根本追不上比赛开场那波尖峰。
Hotstar 表情符号系统的架构实践
原作者:Dedeepya Bonthu,转载自其 Medium 文章
在体育场馆中,观众用欢呼、标语来释放情绪,而电视屏幕前的观众则通过表情符号快速表达自我。当数百万用户同时发送表情时,这就变成了一个技术难题——Hotstar 在互动社交信息流中成功解决了它。
本文从技术视角剖析 Hotstar「Sports Bar」中的表情符号(Emojis)功能:如何实时收集用户信号、压缩为情绪流并动态展示,尤其是在赛事期间面对数十亿次提交的压力。
最初我们使用第三方服务,但性能、稳定性与成本均不理想,最终决定将该核心服务自建。以下为整体架构、关键设计原则及落地影响。
整体架构
架构图略(见原文)。
关键设计原则
可扩展性
系统需支持横向扩展以应对流量增长。通过负载均衡和自动伸缩配置实现资源的动态扩缩容。
分解
系统拆分为多个独立组件,各自承担明确任务,既便于独立扩展,也降低耦合。
异步
异步处理不阻塞资源,从而支持更高并发。这一点将在后文详述。
实现细节
客户端请求处理
客户端通过 HTTP API 提交用户的表情选择。为避免占用连接,API 的重复处理必须离线完成——将数据写入消息队列供下游消费。
消息队列是应用间异步通信的常见机制。对比多种 MQ 后,我们选择了 Kafka:高吞吐、高可用、低延迟,且支持消费组。自运维 Kafka 成本较高,好在 Hotstar 已有基于 Kafka 的数据平台 Knol,完全满足我们的需求。
写入消息队列:异步提升吞吐
同步方式要求等待写入确认才返回成功,适合不可丢失的事务型数据;但表情场景要求极低延迟,偶尔丢失可接受(事实也未出现过)。因此我们选择异步:
- 客户端请求写入本地缓冲,立即返回成功;
- 后台通过 Goroutine 异步将缓冲数据批次刷入 Kafka。
Golang 的并发特性非常契合这一场景。利用 Goroutine 和 Channel,生产者持续消费缓冲数据,按 500ms 间隔或单次 20000 条消息批量刷入 Kafka Broker,兼顾性能与稳定性。
数据处理:Spark 微批聚合
目标是从 Kafka 消费数据流,按小时间窗口计算聚合结果,保证用户感知的实时性。
对比 Flink、Storm、Kafka Streams 后,我们选用 Spark Streaming。它天然支持微批处理和聚合,社区生态也更成熟。Spark 作业每 2 秒计算一次聚合,将结果写入另一个 Kafka 队列。
数据投递:PubSub
投递层使用 Hotstar 自研的 PubSub 实时消息基础设施。我们写了一个简单的 Python Kafka 消费者,按 Spark 批处理周期消费聚合结果,归一化后选出热度最高的表情,推送至 PubSub。客户端订阅后展示动画效果。
落地效果
Emojis 功能在 Hotstar 大获成功。2019 年 ICC 板球世界杯期间,系统从 5583 万用户处收到约 50 亿个表情。截至目前,累计已处理超过 65 亿个表情。
扩展:投票功能
Emojis 与投票本质上都是“近实时处理可量化的用户行为”。我们将该架构复用到投票场景,支撑了印度多个大型真人秀节目(如 Dance Plus、Bigg Boss)的官方投票,累计处理约 30 亿张投票。同一套基础设施还可直接用于投票调查和趣味问答。
原文来源:High Scalability