简介:这是一份面向Java后端与社交应用开发者的抽象库源码,聚焦Timeline模式朋友圈、微博、消息推送、feed流与IM通讯等场景,旨在解决数据流分发、实时推送与消息有序性等工程难题,适合具备一定Java基础、希望快速搭建社交功能骨架的中高级开发者。压缩包共71个文件,约1.59MB,以51个java源码为核心,辅以9个xml配置、3个md说明文档及ftl模板、yml配置等,模块划分清晰,涵盖timeline-core、timeline-store及其redis与memory实现、timeline-mq-starter-redis、timeline-starter以及iot、im等示例工程,便于按需裁剪与二次扩展。已有83人学习下载。读者可从中获得时间线存储与索引、数据流分发、消息推送、feed流聚合及基础IM通信的参考实现,理解多存储后端与消息队列的接入方式,并借助示例工程快速验证与定制业务逻辑。
1. 从朋友圈红点到 IM 会话:Timeline 抽象库到底在解决什么
你打开微信,朋友圈的红点、订阅号的更新、聊天列表的未读,本质上是同一类东西:一条按时间倒序排列、按关系链过滤、按活跃度重排的数据流。Java 后端做社交或通讯类业务,绕不开 Timeline 模式——微博首页、朋友圈、消息推送、feed 流、IM 通讯,底层都是它。但多数团队的做法是每个业务写一套:朋友圈一套 Redis ZSet,推送一套 MQ 消费,IM 又一套会话表。代码重复、逻辑分散、一致性难保证。这个抽象库要干的事,就是把「数据流之间的分发」抽出来:生产者只管写,库负责按 Timeline 模式扇出到不同消费端。适合谁?正在做社交 feed、消息中心、IM 会话列表的 Java 后端,尤其是被多套时间线逻辑折磨过的团队。下面从模型、存储、分发、推送、IM 五个层面拆开讲,每一步都给可复现的代码和参数。
2. Timeline 模型与存储选型:推模式、拉模式还是推拉结合
2.1 三种 Timeline 模式的成本对比
Timeline 模式的核心问题只有一个:用户 A 发了一条内容,怎么让关注 A 的人看到。常见做法有三种。
推模式(fan-out on write):A 发布时,立刻把这条内容写进所有粉丝的收件箱。读的时候直接读自己的收件箱,极快。代价是写放大——一个大 V 有 1000 万粉丝,发一条要写 1000 万次。拉模式(fan-out on read):A 发布只写自己的发件箱,粉丝读的时候去拉所有关注人的发件箱再合并排序。写便宜,读昂贵,关注 2000 人时每次刷新要查 2000 个发件箱。推拉结合:普通用户走推模式,大 V 走拉模式,读的时候把推来的收件箱和大 V 的发件箱合并。
| 模式 | 写成本 | 读成本 | 适用场景 | 典型阈值 |
|---|---|---|---|---|
| 推模式 | 高(粉丝数×) | 低(一次查询) | 普通用户、粉丝少 | 粉丝 < 1 万 |
| 拉模式 | 低(一次写入) | 高(关注数×) | 大 V、明星账号 | 粉丝 > 10 万 |
| 推拉结合 | 中 | 中 | 混合场景 | 按粉丝数分流 |
选型理由:没有银弹。粉丝数中位数低于 500 的产品,直接推模式最简单,别过早优化。只有当出现单账号粉丝过万且发布频繁时,才引入拉模式分流。这个抽象库的价值就在于把「推」和「拉」做成可插拔的策略,业务代码不感知。
2.2 用 Redis ZSet 落地收件箱的最小代码
收件箱用 Redis ZSet 存,score 用时间戳(毫秒),member 用内容 ID。这样按 score 倒序 range 就是时间线。
// TimelineInbox.java public class TimelineInbox { private final JedisPool jedisPool; private static final int MAX_SIZE = 1000; // 每个收件箱最多保留 1000 条 public TimelineInbox(JedisPool jedisPool) { this.jedisPool = jedisPool; } // 写入一条内容到指定用户的收件箱 public void push(long userId, long contentId, long timestamp) { String key = "timeline:inbox:" + userId; try (Jedis jedis = jedisPool.getResource()) { jedis.zadd(key, timestamp, String.valueOf(contentId)); // 裁剪:只保留最新的 MAX_SIZE 条,防止内存无限增长 jedis.zremrangeByRank(key, 0, -(MAX_SIZE + 1)); jedis.expire(key, 7 * 24 * 3600); // 7 天过期,冷用户自动清理 } } // 分页读取收件箱,offset 从 0 开始 public List<Long> page(long userId, int offset, int count) { String key = "timeline:inbox:" + userId; try (Jedis jedis = jedisPool.getResource()) { Set<String> ids = jedis.zrevrange(key, offset, offset + count - 1); return ids.stream().map(Long::parseLong).collect(Collectors.toList()); } } }逻辑说明:zadd的 score 用毫秒时间戳,保证同一毫秒内多条内容按 member 字典序排列,不会丢。zremrangeByRank按排名裁剪,rank 0 是最旧的一条,-(MAX_SIZE+1)表示保留最新 MAX_SIZE 条。expire设 7 天是权衡:活跃用户每次写入会刷新过期时间,不活跃用户 7 天后收件箱自动释放,省内存。
参数说明:MAX_SIZE 设 1000 是经验值,对应大约 20 页(每页 50 条)。超过 1000 条的内容用户几乎不会翻到,保留反而浪费内存。如果你的产品有「查看历史」需求,可以调大到 5000,但要监控 Redis 内存。expire 时间根据产品活跃度调整,日活高的产品可以设 3 天,因为用户每天都会刷新。
提示:ZSet 的 member 必须唯一。如果同一条内容可能被重复推送(比如编辑后重新发布),用 contentId + 版本号做 member,否则 zadd 会覆盖旧 score,导致时间线错乱。
3. 数据流分发:把发布事件扇出到多个 Timeline
3.1 分发器的抽象接口设计
抽象库的核心是「数据流间的分发」。一条内容发布后,要同时进入:发布者的发件箱、粉丝的收件箱、推送队列、IM 会话列表。如果每加一个消费端就改一次发布代码,维护会失控。常见做法是定义 TimelineEvent 和 TimelineSink 两个接口,发布者只发事件,分发器负责路由。
// TimelineEvent.java public class TimelineEvent { private long contentId; private long authorId; private long timestamp; private EventType type; // POST, COMMENT, LIKE, MESSAGE private Map<String, Object> payload; // 扩展字段 // getter/setter 省略 } // TimelineSink.java public interface TimelineSink { String name(); // sink 名称,用于配置路由 boolean supports(EventType type); // 是否处理该类型事件 void onEvent(TimelineEvent event); // 处理逻辑 }逻辑说明:supports让每个 sink 自己声明关心哪些事件类型,分发器只调用返回 true 的 sink。这样新增一个「推送 sink」时,不需要改分发器,只需要注册进去。payload用 Map 存扩展字段,避免每加一个业务字段就改 TimelineEvent 的类结构。
参数说明:EventType枚举建议至少包含 POST(发帖)、COMMENT(评论)、LIKE(点赞)、MESSAGE(私信)。如果你的业务有转发、@提及,继续加枚举值。payload里放渲染需要的字段,比如内容摘要、作者昵称、头像 URL,这样消费端不用回查数据库。
3.2 分发器的注册与异步执行
分发器本身要支持同步和异步两种模式。同步用于强一致场景(比如 IM 消息必须立刻可见),异步用于最终一致场景(比如推送可以延迟几秒)。
// TimelineDispatcher.java public class TimelineDispatcher { private final List<TimelineSink> sinks = new CopyOnWriteArrayList<>(); private final ExecutorService asyncPool = Executors.newFixedThreadPool(8); public void register(TimelineSink sink) { sinks.add(sink); } // 同步分发:所有 sink 执行完才返回 public void dispatchSync(TimelineEvent event) { for (TimelineSink sink : sinks) { if (sink.supports(event.getType())) { sink.onEvent(event); } } } // 异步分发:提交到线程池,不阻塞发布者 public void dispatchAsync(TimelineEvent event) { for (TimelineSink sink : sinks) { if (sink.supports(event.getType())) { asyncPool.submit(() -> { try { sink.onEvent(event); } catch (Exception e) { // 单个 sink 失败不影响其他 sink log.error("sink {} failed for event {}", sink.name(), event.getContentId(), e); } }); } } } }逻辑说明:CopyOnWriteArrayList保证注册时的线程安全,读多写少场景下性能好。dispatchAsync里每个 sink 单独提交任务,一个 sink 抛异常不会影响其他 sink。这是血泪经验:早期版本把所有 sink 放在一个任务里,推送服务超时导致收件箱写入也被拖死。
参数说明:线程池大小 8 是起步值,按 sink 数量和 QPS 调整。如果推送 sink 的 RT 是 200ms,QPS 是 50,那至少需要 10 个线程。建议给每个 sink 配独立的线程池,避免互相影响。dispatchSync只用于 IM 消息这种必须立刻可见的场景,其他一律异步。
注意:异步分发下,事件对象必须是不可变的,或者深拷贝后再提交。否则多个线程同时读同一个 payload,可能读到被修改后的值。
4. 消息推送与 IM 通讯的接入:从事件到端上
4.1 推送 sink 的实现与去重
推送 sink 监听 POST 和 MESSAGE 事件,调用推送网关。这里最大的坑是重复推送:分发器重试、MQ 重投、网络抖动都会导致同一条内容推两次。端上用户看到两条一样的通知,体验极差。
// PushSink.java public class PushSink implements TimelineSink { private final PushGateway gateway; private final RedisTemplate<String, String> redis; @Override public String name() { return "push"; } @Override public boolean supports(EventType type) { return type == EventType.POST || type == EventType.MESSAGE; } @Override public void onEvent(TimelineEvent event) { // 幂等键:contentId + 接收者ID String dedupKey = "push:dedup:" + event.getContentId() + ":" + event.getAuthorId(); Boolean first = redis.opsForValue().setIfAbsent(dedupKey, "1", Duration.ofMinutes(10)); if (Boolean.FALSE.equals(first)) { return; // 已推送过,直接跳过 } gateway.send(event.getAuthorId(), buildPayload(event)); } }逻辑说明:setIfAbsent是 Redis 的原子操作,只有第一次设置成功才返回 true。10 分钟窗口覆盖了绝大多数重试场景。如果推送失败需要重试,重试前要删掉 dedupKey,否则重试会被幂等逻辑挡住。
参数说明:dedupKey 的过期时间设 10 分钟,是因为推送重试通常在秒级到分钟级。如果你的 MQ 重投延迟可能到小时级,调到 1 小时。buildPayload里要控制 payload 大小,APNs 限制 4KB,FCM 限制 4KB,超过会被截断。
4.2 IM 会话列表的 Timeline 化
IM 的会话列表本质也是一个 Timeline:按最后一条消息时间倒序排列。但和 feed 不同的是,同一个会话的新消息要「顶」到最前面,而不是追加。用 Redis ZSet 时,member 用会话 ID,score 用最后消息时间,每次新消息 zadd 更新 score 即可。
// ImSessionTimeline.java public class ImSessionTimeline { private final JedisPool jedisPool; public void updateSession(long userId, long sessionId, long lastMsgTime, String preview) { String key = "im:sessions:" + userId; try (Jedis jedis = jedisPool.getResource()) { // 更新会话的活跃时间,已存在则更新 score jedis.zadd(key, lastMsgTime, String.valueOf(sessionId)); // 单独存预览文本,用于列表渲染 jedis.hset("im:preview:" + userId, String.valueOf(sessionId), preview); jedis.expire(key, 30 * 24 * 3600); } } public List<SessionItem> list(long userId, int count) { String key = "im:sessions:" + userId; try (Jedis jedis = jedisPool.getResource()) { Set<String> sessionIds = jedis.zrevrange(key, 0, count - 1); List<SessionItem> result = new ArrayList<>(); for (String sid : sessionIds) { String preview = jedis.hget("im:preview:" + userId, sid); result.add(new SessionItem(Long.parseLong(sid), preview)); } return result; } } }逻辑说明:zadd对已存在的 member 会更新 score,这正是「顶到最前」的效果。预览文本单独用 Hash 存,因为 ZSet 的 member 只能存一个字符串,塞不下预览。expire设 30 天,IM 会话比 feed 更重要,保留时间长一些。
参数说明:count一般取 50,对应一屏的会话数。如果产品有「置顶会话」,置顶的 sessionId 单独存一个 Set,读取时先取置顶再取普通,合并后返回。预览文本要截断到 50 字符以内,避免 Hash 过大。
提示:IM 会话的未读数不要存在 Timeline 里,单独用计数器(Redis String 的 incr)。Timeline 只负责排序,不负责计数,职责分离。
5. 避坑与排查:Timeline 抽象库落地时的五个翻车现场
5.1 收件箱写入成功但读取为空
现象:发布内容后,粉丝刷新时间线看不到,但 Redis 里确实有数据。原因:收件箱 key 的过期时间被意外刷新或覆盖。比如两个线程同时写同一个用户的收件箱,一个设了 7 天过期,另一个设了 1 秒过期(测试代码残留),后者覆盖前者。解决:过期时间统一在配置中心管理,代码里不允许硬编码。排查时用TTL key看实际过期时间,和预期对比。
5.2 大 V 发布导致 Redis 阻塞
现象:某个粉丝过百万的账号发布内容后,Redis 出现秒级卡顿,其他接口超时。原因:推模式同步写 100 万个收件箱,单次操作耗时过长。解决:大 V 走拉模式,发布只写发件箱;或者推模式改异步,用 MQ 削峰,消费者分批写入。排查时看 Redis 的slowlog,如果出现zadd耗时超过 10ms,就是写放大问题。
5.3 时间线排序错乱
现象:用户看到的时间线里,旧内容排在新内容前面。原因:score 用了秒级时间戳,同一秒内多条内容的 score 相同,ZSet 按 member 字典序排列,导致顺序随机。解决:score 用毫秒时间戳,如果同一毫秒还有并发,在 score 后加一个自增序列(比如timestamp * 1000 + seq % 1000)。排查时用zrange key 0 -1 WITHSCORES看 score 是否单调。
5.4 异步分发丢事件
现象:推送偶尔丢失,日志里没有异常。原因:asyncPool.submit返回的 Future 没有被检查,任务被线程池拒绝时(队列满)抛 RejectedExecutionException,但被吞掉了。解决:给线程池设置合理的拒绝策略,比如 CallerRunsPolicy(调用者线程执行),或者用有界队列 + 自定义拒绝处理器记录日志。排查时监控线程池的getQueue().size()和getRejectedExecutionHandler()的计数。
5.5 IM 消息已读状态不同步
现象:用户在 A 设备读了消息,B 设备的未读数没清零。原因:已读状态只更新了当前设备的本地缓存,没有通过 Timeline 分发到其他设备。解决:已读事件也走分发器,注册一个ReadReceiptSink,收到已读事件后更新服务端的未读计数,并通过长连接推给其他设备。排查时检查已读事件的EventType是否被正确路由。
6. 进阶技巧:用滑动窗口做 Timeline 的冷热分离
Timeline 数据有个特点:越新的内容访问越频繁,越旧的内容几乎没人看。全量放 Redis 成本高,全量放 MySQL 查询慢。我一般用滑动窗口做冷热分离:热数据(最近 7 天)放 Redis ZSet,冷数据(7 天前)归档到 MySQL,读取时先查 Redis,不够再查 MySQL 并回填。
// HybridTimeline.java public class HybridTimeline { private final TimelineInbox hot; // Redis private final ColdStorage cold; // MySQL public List<Long> page(long userId, int offset, int count) { List<Long> result = hot.page(userId, offset, count); if (result.size() < count) { // 热数据不够,从冷存储补 int need = count - result.size(); List<Long> coldData = cold.query(userId, offset + result.size(), need); result.addAll(coldData); // 回填到 Redis,下次直接命中 for (Long id : coldData) { hot.push(userId, id, getTimestamp(id)); } } return result; } }逻辑说明:先查热数据,不够再查冷数据,查到后回填 Redis。回填时要注意 score 用原始时间戳,不能用当前时间,否则冷数据会排到热数据前面。getTimestamp(id)从内容元数据里取发布时间。
参数说明:热数据窗口 7 天是经验值,按产品活跃度调整。日活高的产品可以缩到 3 天,省内存。冷存储用 MySQL 时,按userId分表,每张表加(userId, publish_time)联合索引,查询走覆盖索引。回填的批量大小控制在 100 条以内,避免单次 Redis 写入过大。
注意:回填会导致冷数据被重复写入 Redis,如果多个请求同时回填同一条数据,zadd 是幂等的(相同 member 更新 score),不会产生重复。但要注意并发回填时的 Redis 连接数,建议用 pipeline 批量写。
这套方案我在两个项目里跑过,最大的教训是:别一开始就上推拉结合和冷热分离。先用最简单的推模式 + Redis ZSet 跑通,等真的遇到大 V 瓶颈或内存瓶颈再优化。过早抽象比不抽象更可怕,因为你会花大量时间维护一套没人用的策略框架。希望帮到你。
本文还有配套的精品资源,点击获取