资讯动态

x-algorithm 如何实现亚毫秒内网取数:Thunder 基于 Kafka 与内存 PostStore 的实时帖子存储设计

发布时间:2026/10/9 15:41:01 来源:尧图企业网站定制
x-algorithm 如何实现亚毫秒内网取数Thunder 基于 Kafka 与内存 PostStore 的实时帖子存储设计【免费下载链接】x-algorithmAlgorithm powering the For You feed on X项目地址: https://gitcode.com/gh_mirrors/xa/x-algorithmx-algorithm 是驱动 X「For You」信息流推荐算法的开源项目。其中 thunder/ 模块以「Kafka 事件流 纯内存 PostStore」的组合实现了亚毫秒级的内网In-Network帖子取数你关注账号发新帖的瞬间帖子就会落进内存索引信息流请求到来时直接从内存中秒级返回无需触碰任何磁盘数据库。本文将带你完整拆解这套实时帖子存储设计数据如何从 Kafka 流入内存、五张哈希表如何组织、一次查询如何做到亚毫秒、过期数据又如何被自动裁剪。一、Thunder 在整个推荐链路中的位置x-algorithm 的 For You 信息流由两大来源组装来源模块作用内网In-Networkthunder/内存中保留你关注账号的近期帖子外网Out-of-Networkphoenix/、simclusters/找到你未关注账号的帖子两者合并后由同一模型统一排序。Thunder 负责内网这一半它对外暴露 gRPC 接口GetInNetworkPosts由信息流混排服务通过 home-mixer/sources/thunder_source.rs 并行调用。 为什么需要内网这条独立通路——你关注的人刚发的帖子是信息流中最新鲜、点击意愿最强的候选必须优先且极快地取到。二、数据入口Kafka 多分区并行消费帖子数据的入口是 Kafka 事件流。服务启动时thunder/main.rs 会先启动 Kafka 消费者并等待追平积压catchup后才将服务标记为就绪——这保证对外提供查询时内存里已经有足够新鲜的数据。1. 分区与线程的切分消费逻辑在 thunder/kafka/tweet_events_listener_v2.rs默认32 个处理线程瓜分32 个 Kafka 分区参数见 thunder/args.rs每个线程独占一部分分区互不干扰地并行拉取消息。每个分区还挂着一个 lag 监控任务持续上报各分区消费延迟方便排查数据新鲜度。2. 批量拉取与批量写入线程循环中按批次默认kafka_batch_size 1000条拉取消息攒满一批后反序列化为轻量的LightPost只保留帖子 ID、作者、时间、回复/转写目标等十几个字段见 thunder/deserializer.rs将创建事件批量写入 PostStore删除事件批量标记失效insert_posts/mark_as_deleted追平阶段用信号量控制并发避免写入线程挤占 CPU。⚡ 关键细节追平时每个线程会检测本组分区的总 lag当 lag 小于「分区数 × 批大小」时才向主流程报告该线程追平——所有线程都追平后服务才置为 readythunder/main.rs。此外还有一条 v1 消费链路 thunder/kafka/tweet_events_listener.rs它消费原始tweet_events话题解析后重新生产精简的innetwork_post事件到 Phoenix 集群的 Kafka供 serving 集群消费——即先清洗、再分发的二级管道配置见 thunder/kafka_utils.rs。三、内存 PostStore五张 DashMap 组成的实时索引核心数据结构在 thunder/posts/post_store.rs整个存储只有五张DashMap分片并发哈希表结构键 → 值用途posts帖子 ID →ArcCompactPost帖子主体压缩成 12 个定长字段original_posts_by_user用户 ID →VecDequeTinyPost该用户的原创帖按时间排列secondary_posts_by_user用户 ID →VecDequeTinyPost该用户的回复/转推帖video_posts_by_user用户 ID →VecDequeTinyPost该用户的视频帖含转推的源视频帖deleted_posts帖子 ID →true墓碑表防止已删帖被迟到消息复活三个省内存的小心机TinyPost 只有 16 字节帖子 ID 时间戳L20-L24用户的帖子清单里不存正文只存指针真正的帖子体统一放在posts表中按需回查CompactPost 是定长结构L35-L48用0代替Option占位比协议对象更紧凑配合Arc做到零拷贝共享每个用户的清单有硬上限MAX_POSTING_LIST_SIZE 5000thunder/config.rs写入时若超限自动弹出队首最老的帖子并从posts表中同步移除内存占用因此被严格封顶。删除与复活防护帖子删除时走mark_as_deletedL111-L125从posts移除同时把 ID 记入墓碑表deleted_posts。插入路径上有一行关键判断——若帖子已在墓碑表中则直接跳过。这就防止了 Kafka 乱序时一条迟到的创建事件把已删除的帖子重新塞回内存。四、查询路径一次取数如何做到亚毫秒当 home-mixer 的 ThunderSource 带着「关注列表 已看过帖子 ID」发起 gRPC 请求后thunder/thunder_service.rs 的处理路径全部在内存中完成限流保护全局 QPS 限流器默认 8000 QPS先检查配额超了直接返回resource_exhausted防止突发流量打穿服务内存查询请求被丢进阻塞线程执行get_all_posts_by_usersthunder/posts/post_store.rs——对每个关注账号从其VecDeque尾部最新反向遍历套用每作者上限原创帖默认 50 条、回复转推默认 30 条thunder/config.rs超时熔断遍历中持续检查request_timeout一旦超时立即中止并计数保证最坏情况下延迟也有上界新鲜度排序最后按created_at倒序取前 N 条默认 1200 条返回thunder/thunder_service.rs。整个路径没有任何网络 IO 和磁盘 IO——两次哈希查找 顺序扫描定长结构体这正是亚毫秒的来源。五、容量与新鲜度的长期治理内存存储最大的风险是越攒越多。Thunder 用三重机制保证数据始终新鲜且有限保留期retention默认172800秒2 天thunder/args.rs。新帖子入库前就检查时效超龄消息直接丢弃自动裁剪服务启动后每2 分钟运行一次trim_old_poststhunder/main.rs → thunder/posts/post_store.rs。裁剪时按用户ID % 5分桶每轮只处理 1/5 的用户把开销摊平到 10 分钟内避免周期性卡顿启动时预热排序finalize_init会在追平 Kafka 后对全部用户清单按时间重排并分 5 轮执行首轮裁剪L142-L154保证从可服务那一刻起索引就是有序的。此外每 5 秒的统计任务会持续上报用户数、帖子总量、各清单容量等指标POST_STORE_*系列thunder/metrics.rs让运维可以实时观察内存水位与数据新鲜度。六、设计要点总结设计点做法效果数据新鲜度Kafka 32 分区并行消费 追平后才 ready新帖秒级可查查询延迟纯内存 DashMap 定长结构体亚毫秒级返回内存上限每用户 5000 条清单 2 天保留期 分桶自动裁剪内存占用可控数据正确性墓碑表防复活、回复归属过滤、每作者限流结果干净、公平服务稳定QPS 限流、单请求超时熔断、分区 lag 监控延迟有上界Thunder 用最朴素的两件东西——消息队列 内存哈希表——解决了推荐系统里最苛刻的新鲜度 vs 延迟矛盾。想深入阅读建议从 thunder/posts/post_store.rs 的五张表定义开始再顺 thunder/kafka/tweet_events_listener_v2.rs 看数据入口最后对照 home-mixer/sources/thunder_source.rs 看它如何被信息流调用即可完整掌握这条亚毫秒取数链路。【免费下载链接】x-algorithmAlgorithm powering the For You feed on X项目地址: https://gitcode.com/gh_mirrors/xa/x-algorithm创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑