实时数据平面
晚高峰客户端行为数千万 IOPS。这个量级不能让 Feed 服务同步写十个下游。总线先吃下来,再按消费者能力扩。
用户
App
BMQ
Flink
特征 KV
Monolith
划走
▶
上报事件
▶
消费
▶
更新短窗特征
▶
样本
▶
更新 embedding
▶
下次 Feed 就能用上
BMQ / Kafka
早期 Kafka 本地盘 + ISR,扩缩容要搬副本,热点 partition 打满单盘。BMQ(SoCC 2024,内部也称 ByteMQ)计算存储分离:Proxy 聚合 Produce,Broker 写分布式文件系统,Consume 可走多层缓存直接读存储。Partition 切 Segment 打散到不同磁盘,热点被摊开。Controller 不再管 ISR,故障可秒级切。公开:约 99.76% Kafka 集群迁走,资源成本约降 70%,单集群 TB/s 吞吐口径。
Flink → KV
消费行为 Topic,按 uid/item 更新窗口计数、序列、实时 CTR。状态在 RocksDB StateBackend;热 key 靠 State Cache 避免 LSM 反复 compaction。聚合结果写入 Abase,给精排点查。同时分流给反作弊、热点检测、计数服务。
Hudi / Hive / Spark
流同时落湖。Hudi 承接样本插入/更新/删除和特征回溯,以及把在线 LSM 存储的 CDC 打到离线可扫的格式。离线训练出稠密网络初始权重和新 item embedding,再导入 Monolith batch 阶段。
闭环时序
- T+0 端上组包上报,接入写入 BMQ。
- T+秒级 Flink 更新短窗特征(刚看过的类目、即时负反馈)。
- T+分钟 Monolith 把被碰到的 embedding 同步到 serving PS。
- T+小时/天 离线重训稠密层、重建 ANN 索引、刷新协同表。
短窗决定「现在不想看」,长窗决定「你是谁」。两条必须并存。