ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Timeline抽象库设计:一招统一朋友圈、微博、推送与IM的数据流

Timeline抽象库设计:一招统一朋友圈、微博、推送与IM的数据流 简介在社交、社区和即时通讯后端开发中时间线Timeline是承载数据流的核心数据结构。无论是朋友圈的好友动态、微博的关注Feed、消息中心的系统通知还是IM的会话消息其本质都是按时间排序的事件流分发与消费。理解推模型Fanout On Write与拉模型Fanout On Read的适用边界是设计高并发读写系统的关键。推模型读路径轻但写放大严重拉模型写放大为零但读时合并代价高实际业务常采用推拉结合策略。通过抽象Timeline事件、游标分页、幂等去重及可插拔分发策略可以将不同业务场景统一到同一套数据流建模体系中。本文从基础概念出发剖析Timeline抽象库如何覆盖Feed流、通知推送与IM通讯的共性帮助后端工程师设计出可扩展、可运维的时序数据流系统。 这几年做社交、社区、内容类后端项目我几乎每隔一段时间就要在一个新项目里把时间线逻辑重新写一遍。朋友圈点了发布粉丝要立刻拉到动态微博大V发一条内容百万级粉丝刷新要看得到后台做消息推送需要按标签把通知推到一批用户收件箱IM要保证两个人之间的消息有序到达且不丢不重。这四类业务看起来天差地别但落到代码层面其实是同一个问题数据流的分发。我做的这个Java抽象库就是把Timeline模式里“数据流如何建模、如何分发、如何被消费”的部分抽出来让业务层只需要关心自己的关系链和业务语义。这篇文章我会直接把整个抽象库的设计思路、核心接口、分发机制以及它如何同时适配朋友圈、微博、消息推送、IM通讯这四类场景讲清楚。代码层面会给出一个最小可用实现并把我在真实业务里踩过的坑一并说出来。适合正在设计feed流、消息中心、通知系统或者单纯想理解时序数据流如何抽象的Java后端工程师阅读。1. Timeline 模式为什么要做一个“抽象库”而不是直接写业务1.1 每家企业都在重复造同一个轮子我见过不少项目组做动态feed流的方式产品经理说要做“关注页”后端就在用户表旁边加一张follow关系表再写一条SQL把关注对象的最新内容查出来然后按时间排序。第一版能跑等数据量上来就开始出问题SQL越写越重、缓存不知道怎么设计、分页越翻越慢。这时候才会有人意识到“关注页”背后其实是一套独立于业务的时间线系统。同样的剧情会在做通知中心、私信模块、广播消息时再演一遍。区别只是表名从follow变成tag、从friend变成conversation排序逻辑和分发逻辑几乎一模一样。我做这个抽象库的初衷很直接把这些场景里公共的部分沉淀下来让业务代码只表达“谁产生了什么数据、应该流向哪些接收者”而不去关心接收者手里的时间线是怎么构建和持久化的。1.2 抽象库的边界在哪里很多工程师一听“抽象库”就往框架方向想但这套东西不是Spring Boot Starter也不是要接管全公司的数据存储。它的边界非常克制不管业务关系链怎么存关注关系在MySQL还是在图数据库跟这个库无关。不管业务事件长什么样朋友圈的图文、IM的文本消息、推送的告警通知都是payload。管的是Timeline本体一个接收者手里那条有序数据流如何追加、如何读取、如何游标分页。管的是数据流之间的分发一条业务数据从生产者产生后按照什么规则被复制到一批接收者的Timeline里。这个边界很重要。如果库管太多业务方会因为它不够灵活而放弃如果管太少又退化成一张表没有抽象价值。Timeline抽象库的核心价值就两个词建模和分发。1.3 这个库解决了什么具体问题第一个问题是语义统一。业务方不再需要各自定义一套“查询好友最新动态”的接口而是统一理解成“读取某条Timeline的某个区间”。第二个问题是写入链路的复用。无论数据来自发帖、发消息、发通知统一走一个Dispatcher入口分发的规则由业务方传入但分发的过程、幂等、重试、游标推进由库来处理。第三个问题是让容量评估变得可计算。一旦业务方理解了自己用的是推模型还是拉模型就能算出“一次发布会产生多少份写放大”这在架构设计阶段非常有用。2. 朋友圈、微博、推送、IM 的底层共性2.1 四种场景的形态对比先把四种典型场景放在一张表里看它们的表象差异会非常清楚场景数据生产者接收者范围时间线维度实时性要求主要读写特点微博/关注feed被关注的作者粉丝集合用户聚合流秒级读多写多大V写放大极端朋友圈好友好友集合用户聚合流秒级好友数有上限关系闭合消息推送系统/运营标签、用户组用户通知流分钟级可接受批量写入为主IM通讯对话参与者会话内成员会话消息流毫秒级严格有序实时性要求高表里的“时间线维度”很关键。微博和朋友圈的用户时间线是由多个生产者聚合而成的消息推送的时间线是由系统单方面写入的IM的时间线则是按会话维度存储的会话双方共享同一条流。2.2 被业务表象掩盖的三个共同点如果把表里的差异剥掉剩下三个在所有场景里都成立的事实每个接收者最终看到的都是一个按时间排序的列表。微博用户看到的是关注对象的动态列表用户看到的是好友动态列表手机用户看到的是通知列表聊天窗口里看到的是消息列表。列表的构建都逃不开“发件箱/收件箱”模型。生产者先产生一条数据放进自己的发件箱再通过某种机制被复制或关联到接收者的收件箱。区别只是复制是写入时做还是读取时做。读取端都需要稳定的游标分页。用户下拉刷新、上滑翻页本质都是在时间线上移动游标。我在设计这个库时就是围绕这三条共性来建模的。业务场景的差异通过“分发策略”和“时间线类型”两个扩展点来消化而不是靠堆接口。2.3 产品语义差异怎么映射到统一模型这里最容易走偏的做法是为了统一强行让IM和微博共用同一套代码路径结果两边都别扭。我的做法是保留“事件”和“分发策略”的抽象允许业务定义自己的category和DispatchStrategy不同的策略对应不同的数据流向。产品层面的“朋友圈好友动态”、“微博关注流”、“通知中心”、“聊天会话”都是Timeline的实例只是它们的生产者、订阅关系、分发时机不同。用户在UI上看到的是不同的页面在系统里看到的是同一种数据结构。3. 推模型、拉模型时间线系统绕不开的路线之争3.1 推模型Fanout On Write的适用边界推模型的核心思路是生产者在写入时就把数据复制到所有接收者的收件箱里。以朋友圈为例你发一条动态系统立刻把它写到每个好友的收件箱时间线里好友刷新时直接读自己的收件箱不需要实时聚合。推模型的优势是读路径极轻。一个用户刷新自己的时间线只需要按游标顺序从自己的收件箱取数据不管关注了多少人读取成本都只和自己的数据量相关。适合朋友圈这种关系链闭合、每个用户好友数有限通常几千以内的场景。但推模型有一个绕不开的天花板写放大。假设一个用户有500个好友他发一条动态要写500份如果一个千万级粉丝的大V也走推模型发一条内容就要写千万份这在任何存储系统里都不可接受。所以朋友圈能用推模型微博不能完全用推模型。3.2 拉模型Fanout On Read和延迟合并拉模型的核心思路是生产者只写入自己的发件箱接收者读取时再实时聚合自己订阅的多个发件箱。微博的关注feed流就是典型大V发一条微博只写一次粉丝刷新时去读取自己关注的作者列表再拉取每个作者的最新内容合并排序。拉模型的优点是写放大为零代价是读放大。一个用户关注了300个作者每次刷新都要拉300个发件箱的最新内容再在内存里合并排序读取延迟和发件箱数量成正比。为了缓解这个问题业界一般会加缓存、并行拉取、做多级合并。拉模型还有一个隐含问题时间线牺牲了绝对的有序性。因为数据是读取时实时归并的如果在聚合过程中某个作者又发了新内容用户可能在同一页里看到“新一条”插在中间导致游标分页变得复杂。3.3 抽象层必须同时支持两者很多自研时间线系统容易犯的错误是选了一种模型就写死后面想切换只能重写。我在抽象库中把“分发策略”作为接口暴露出来推模型和拉模型是两种可插拔策略推模型对应写入时立即fanout到目标收件箱拉模型对应写入时只写生产者发件箱读取时由TimelineReader触发实时聚合。实际业务往往是推拉结合的。大V走拉模型普通用户走推模型热点事件可以用“预聚合”的方式定期把热门作者的动态批量分发到粉丝收件箱IM则因为数据量小但实时性要求高通常直接走推模型。抽象层不替业务做决策它只负责提供两种能力并把“分发后数据一致”的保障做好。具体哪条流用哪种模型由业务在创建Timeline时通过配置声明。4. 核心抽象设计从接口到分发机制4.1 一页纸讲清整体模型这套库的逻辑可以用一句话概括业务事件从生产者进入DispatcherDispatcher根据订阅/策略把事件路由到一批目标Timeline上使用者通过Reader从Timeline里按游标读取有序事件。结合发件箱/收件箱模型看每个生产者也拥有一条“发件箱Timeline”记录他产生的所有事件每个接收者拥有一条“收件箱Timeline”记录他需要消费的事件Dispatcher负责把发件箱的事件按规则复制到接收者的收件箱。读路径上不需要考虑数据从哪里来只需要面对自己的收件箱。4.2 最小接口集定义我设计的核心接口只有四个这也是整个库的心脏public interface Timeline { String timelineKey(); void append(TimelineEvent event); TimelinePage read(Cursor cursor, int limit); } public interface TimelineEvent { String eventId(); long timestamp(); String producerKey(); int category(); byte[] payload(); } public interface Dispatcher { void publish(TimelineEvent event, DispatchStrategy strategy); } public interface TimelineReader { TimelinePage read(String timelineKey, Cursor cursor, int limit); }Timeline接口定义一条时间线的最小行为追加事件和读取事件。TimelineEvent是所有业务载荷的通用包装eventId用于全局幂等timestamp用于排序producerKey标明来源category用于区分业务类型payload携带具体业务内容。我没有把“删除”和“修改”放进接口里因为时序数据流的核心语义是append-only。业务层面的“删动态”“撤回消息”应该通过追加一条“删除标记事件”这种方式表达而不是物理删除历史事件。这个设计决策帮我避开了很多分布式场景下的数据一致性问题。4.3 分发规则的表达DispatchStrategy分发策略是整个库扩展性最强的地方。它本质上是一个“给定一个事件计算出需要投递到哪些Timeline”的函数。我用接口表达public interface DispatchStrategy { ListString resolveTimelineKeys(TimelineEvent event); }推模型就是返回一批接收者收件箱的key拉模型就是返回空列表因为不需要复制。但为了支持拉模式聚合我还需要另一种能力读取一个Timeline时除了自身存储的事件还需要合并其他Timeline的事件。所以TimelineReader有另一个聚合方法public interface TimelineReader { // 聚合读取读取主timeline并合并sources里的事件 TimelinePage mergeRead(String timelineKey, ListString sourceTimelineKeys, Cursor cursor, int limit); }这样微博场景的读取逻辑就变成用户的收件箱里可能只有普通用户主动推送来的事件大V的事件不提前推送而是读取时通过mergeRead动态合并大V的发件箱。4.4 Cursor 游标为什么不是数字而是对象早期我做时间线分页直接用pageNum/pageSize或者lastId都遇到过问题。lastId如果用的是数据库自增id一旦数据迁移或者多库合并顺序就不可靠用offset深翻页则性能越来越差。这个库里的Cursor设计成一个包含位置信息的对象public class Cursor { private final long timestamp; private final String lastEventId; private final boolean forward; public static Cursor start() {...} public static Cursor from(long timestamp, String lastEventId) {...} }timestamp定位时间位置lastEventId防止同一毫秒内有多条事件时出现跳过或重复forward表示向后翻页还是向前拉新。读取时先按timestamp过滤再处理相同时间戳的事件用eventId做排序和去重。这套游标模型在IM场景里同样成立配合eventId的全局唯一性可以做到不丢不重。4.5 幂等和去重分发链路上的基石分发本质上是一个“复制”过程。只要涉及复制就一定会遇到重复网络重试、消费者重放、生产者重发都可能导致同一条事件被追加两次。我在Timeline.append()实现里强制按eventId做唯一性约束。在内存实现里这个是ConcurrentHashMap加TreeSet的组合在Redis实现里用SETNX保证同一个timelineKey eventId只写入一次在MySQL实现里timeline_key event_id建唯一索引。这套幂等逻辑是跨所有存储实现公用的业务方不需要关心重试问题只管把事件交给Dispatcher就行。5. 用这套抽象把四个场景各接一遍5.1 模拟微博关注流拉模型为主推模型为辅// 大V发布一条微博走拉模型只写大V自己的发件箱 String bigV user-10001; TimelineEvent event SimpleEvent.builder() .eventId(UUID.randomUUID().toString()) .timestamp(System.currentTimeMillis()) .producerKey(bigV) .category(Category.FEED) .payload(json.getBytes(StandardCharsets.UTF_8)) .build(); dispatcher.publish(event, event - Collections.emptyList());普通粉丝读取关注页时// 粉丝user-20001读取时间线先读自己的收件箱再merge大V的发件箱 ListString bigVOutboxes followService.getFollowedBigVOutboxKeys(user-20001); TimelinePage page reader.mergeRead( user-20001, // 自己收件箱 bigVOutboxes, // 大V发件箱列表 cursor, 20 );普通用户发布时走的还是推模型DispatchStrategy返回所有粉丝的收件箱key。只有大V才切换策略避免写放大。这样一套接口两种策略就都接上了。5.2 模拟朋友圈典型推模型朋友圈的场景是最契合推模型的关系链闭合、好友数有上限。假设用户user-30001发了一条动态他的好友列表直接从关系服务查出来ListString friendIds relationService.friendIdsOf(user-30001); dispatcher.publish(event, evt - friendIds.stream() .map(fid - inbox: fid) .collect(Collectors.toList()));每条动态都被复制到所有好友的收件箱。好友刷新时读取inbox:{uid}天然有序不需要任何实时合并。这也是朋友圈产品体验流畅的原因之一读路径轻到几乎没有计算量。朋友圈场景还有一个独有的需求谁可以看。这类“可见性过滤”不适合放在时间线查询链路里做实时过滤因为会拖慢读路径。我的做法是把可见性条件编码进分发阶段如果一条动态是“仅部分好友可见”那么DispatchStrategy返回的收件箱列表直接排除不可见好友。过滤提前到写入路径读取端完全无感知。5.3 模拟消息推送按标签批量分发消息推送和社交feed最大的区别在于接收者集合不是“关系链”而是“标签或分组”。运营选一个标签系统把通知发给这个标签下的所有用户。String tagId tag:promotion-2024; ListString userIds tagService.userIdsByTag(tagId); dispatcher.publish(event, evt - userIds.stream() .map(uid - notice: uid) .collect(Collectors.toList()));这个场景的写放大规模通常是百万级。所以推送场景下我更推荐异步分发Dispatcher只把事件写入一个待分发队列比如用内存队列或消息队列真正的fanout由后台worker批量执行。抽象库的Dispatcher接口本身就是异步实现和同步实现可以替换的业务方按吞吐量要求选即可。推送时间线还需要处理“过期失效”问题。比如一条优惠券通知活动结束后再展示没有意义。我一般会追加一条“过期标记事件”而不是物理删除通知客户端收到标记事件后做本地隐藏。这样可以保留完整历史也避免物理删除在分库分表场景下的麻烦。5.4 模拟IM通讯会话维度的消息流IM场景和前面三个都不一样它不强调“一对多广播”而是“一对一会话双方共享一条有序流”。所以IM的时间线key不是用户维度而是会话维度。String conversationKey conv:user-40001:user-40002; TimelineEvent event SimpleEvent.builder() .eventId(snowflake.nextIdStr()) .timestamp(System.currentTimeMillis()) .producerKey(user-40001) .category(Category.IM) .payload(在吗.getBytes(StandardCharsets.UTF_8)) .build(); dispatcher.publish(event, evt - List.of(conversationKey));会话双方读取时都读conv:user-40001:user-40002这一条时间线通过游标做增量拉取。IM场景对顺序要求非常严格所以timestamp应该由服务端生成不能信任客户端时间eventId用雪花算法生成保证全局唯一且趋势递增。IM的幂等逻辑比社交场景更重要。用户弱网重试、客户端重发消息同一个eventId会被追加多次Timeline层的唯一索引会兜住重复写入客户端只需要按eventId做去重即可。这也是我把eventId设计成接口必填字段的原因。5.5 一个内存版最小实现为了让这套抽象不悬空我写了一个内存版实现逻辑足够简单适合二次开发和理解核心流程public class InMemoryTimeline implements Timeline { private final String key; private final NavigableMapLong, ListTimelineEvent events new TreeMap(); public InMemoryTimeline(String key) { this.key key; } Override public String timelineKey() { return key; } Override public synchronized void append(TimelineEvent event) { events.computeIfAbsent(event.timestamp(), k - new ArrayList()) .add(event); } Override public synchronized TimelinePage read(Cursor cursor, int limit) { // 游标过滤 按 eventId 排序 截断 limit ListTimelineEvent result events.tailMap(cursor.timestamp(), false) .entrySet().stream() .flatMap(e - e.getValue().stream()) .filter(e - isAfter(e, cursor)) .sorted(Comparator.comparing(TimelineEvent::timestamp) .thenComparing(TimelineEvent::eventId)) .limit(limit) .collect(Collectors.toList()); if (result.isEmpty()) { return TimelinePage.empty(); } TimelineEvent last result.get(result.size() - 1); return new TimelinePage(result, Cursor.from(last.timestamp(), last.eventId())); } }这段代码里最容易被忽略的是tailMap(cursor.timestamp(), false)的边界条件。游标记录的是“上一次读取的最后一条事件”下回读取要从这个时间点的后面开始所以用false表示不包含当前时间戳。相同timestamp的多条事件则通过eventId字符串排序保证全局顺序稳定。6. 从抽象到落地缓存、分页、热点这些坎儿6.1 存储选型不能一套打天下很多人拿到这个库首先问Timeline数据应该存哪里我的经验是分场景社交feed流推荐Redis的Sorted SetZADD按时间戳写入ZREVRANGE按游标读取天然支持按分值范围分页。数据量大可以做冷热分层旧数据下沉到MySQL或者对象存储。IM会话流推荐用类Cassandra的宽表存储或者直接用成熟的IM存储保证多端同步的可靠性。内存库结合消息队列也可以做前置层。消息推送推送的读频率远低于写频率但一次性写入量大适合用MySQL批量插入加Redis缓存热点数据。存储实现和抽象接口是解耦的。我始终把Timeline、Dispatcher、TimelineReader作为接口存储细节全部放在实现类背后。这样业务在早期用内存实现验证模型数据量上来之后无缝切换到Redis或数据库实现不需要改业务代码。6.2 时间线分页游标比Offset可靠得多用offset做时间线分页的问题其实不只是性能。你在offset翻页的时候如果前面插入了新数据整页内容会整体后移用户会看到重复或不连续的内容。游标分页不会受这个影响因为它锚定的是“上一次读到的位置”新数据来了只会出现在游标之后。IM场景的增量同步更是离不开游标。客户端每隔几秒拉一次新消息带上的就是上次同步的游标。服务端只需要返回游标之后的事件天然做到“只拉增量”。我在实际使用中养的的习惯是把游标序列化成不透明字符串下发给客户端。客户端不解析、不修改每次原样返回。这样服务端可以自由演进游标内部结构不用考虑兼容性。6.3 大V热点的缓解手段拉模型虽然解决了大V的写放大问题但把压力转移到了读路径百万粉丝同时刷新大V的发件箱会被反复读取这个key会成为热点。我常用的缓解手段有三个本地缓存在应用层加大V发件箱的短时本地缓存比如几百毫秒到一秒能挡住大部分重复请求。时间片拆分把大V的发件箱按小时拆成多个子时间线读取时并行拉取。这样单个key的压力就分散了。多级合并缓存热门关注的“聚合结果”可以预计算按秒级更新粉丝直接读预聚合结果而不是实时合并。这些手段都符合抽象库的接入方式缓存逻辑写在TimelineReader实现里业务层完全无感。这也是接口抽象带来的额外好处调优可以发生在框架内部不必层层传递。6.4 时间戳统一用服务端时间别信客户端时间线排序最怕的就是时间错乱。客户端本地时间可能被用户修改也可能因为时区问题产生偏差。如果服务端容忍客户端时间戳直接参与排序就会出现“新消息排在旧消息后面”这种严重事故。我的规则是所有进入Timeline的事件的timestamp必须在服务端生成客户端传的时间只作为业务字段保存不参与时间线排序。IM场景尤其要注意这一点服务端收到消息后立即打点再交给Dispatcher分发。6.5 关于持久化、备份和迁移的提醒内存版实现适合做验证不能直接上生产。接MySQL实现时建议timeline_key、event_id、timestamp建联合索引或唯一索引接Redis实现时要注意RDB和AOF的配置因为纯内存版在宕机时会丢数据。IM场景如果允许少量丢失可以做异步刷盘但如果要求严格不丢需要引入可靠消息队列和多副本存储。我见过一个项目因为没给timeline数据做备份误删一条大V发件箱记录后几十万粉丝的时间线同时缺了那条内容排查了整整一下午。从那以后我所有Timeline实现类都会默认加上“事件追加审计日志”写一条真实数据之前先写一条审计记录万一出问题可以靠审计日志重放恢复。7. 一次实际集成把我坑醒的教训最后说一个我自己犯过的错误算是给这个抽象库做一次“现实检验”。之前把一个消息中心业务迁到这套模型上初期只考虑了“写入、读取、游标”三个动作没想过“大V发件箱被粉丝并发拉取”的压力。上线后运营做了一次全量推送目标用户300万推送事件全部走了推模型瞬间把存储写入打满数据库慢查询暴涨。后来我把策略改成运营推送走异步批量分发同时普通用户走推、大V类账号走拉问题才缓解。另一件事是eventId的生成规范。最开始有的业务方用时间戳加随机数拼eventId并发一高就出现重复幂等校验直接把合法事件挡在外面。后来统一要求必须用雪花算法或者UUID每个事件的eventId全局唯一幂等逻辑才真正可靠。根据我个人的经验这套Timeline抽象最大的价值不是省了那几张表的代码而是逼着业务方在写第一行代码之前就把“数据从哪来、流向谁、怎么排序、怎么分页”这四个问题想清楚。越早理清数据流后期越少返工。如果你也正在设计feed流或者消息系统建议先画一张数据流图再套这套抽象能少踩不少我踩过的坑。本文还有配套的精品资源点击获取
返回列表