ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Java后端Timeline抽象库:从朋友圈到IM的推拉模式与Redis ZSet实践

Java后端Timeline抽象库:从朋友圈到IM的推拉模式与Redis ZSet实践 简介这是一份面向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 writeA 发布时立刻把这条内容写进所有粉丝的收件箱。读的时候直接读自己的收件箱极快。代价是写放大——一个大 V 有 1000 万粉丝发一条要写 1000 万次。拉模式fan-out on readA 发布只写自己的发件箱粉丝读的时候去拉所有关注人的发件箱再合并排序。写便宜读昂贵关注 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 ListLong page(long userId, int offset, int count) { String key timeline:inbox: userId; try (Jedis jedis jedisPool.getResource()) { SetString ids jedis.zrevrange(key, offset, offset count - 1); return ids.stream().map(Long::parseLong).collect(Collectors.toList()); } } }逻辑说明zadd的 score 用毫秒时间戳保证同一毫秒内多条内容按 member 字典序排列不会丢。zremrangeByRank按排名裁剪rank 0 是最旧的一条-(MAX_SIZE1)表示保留最新 MAX_SIZE 条。expire设 7 天是权衡活跃用户每次写入会刷新过期时间不活跃用户 7 天后收件箱自动释放省内存。参数说明MAX_SIZE 设 1000 是经验值对应大约 20 页每页 50 条。超过 1000 条的内容用户几乎不会翻到保留反而浪费内存。如果你的产品有「查看历史」需求可以调大到 5000但要监控 Redis 内存。expire 时间根据产品活跃度调整日活高的产品可以设 3 天因为用户每天都会刷新。提示ZSet 的 member 必须唯一。如果同一条内容可能被重复推送比如编辑后重新发布用 contentId 版本号做 member否则 zadd 会覆盖旧 score导致时间线错乱。3. 数据流分发把发布事件扇出到多个 Timeline3.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 MapString, 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 ListTimelineSink 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 是 200msQPS 是 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 RedisTemplateString, 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 限制 4KBFCM 限制 4KB超过会被截断。4.2 IM 会话列表的 Timeline 化IM 的会话列表本质也是一个 Timeline按最后一条消息时间倒序排列。但和 feed 不同的是同一个会话的新消息要「顶」到最前面而不是追加。用 Redis ZSet 时member 用会话 IDscore 用最后消息时间每次新消息 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 ListSessionItem list(long userId, int count) { String key im:sessions: userId; try (Jedis jedis jedisPool.getResource()) { SetString sessionIds jedis.zrevrange(key, 0, count - 1); ListSessionItem 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 ListLong page(long userId, int offset, int count) { ListLong result hot.page(userId, offset, count); if (result.size() count) { // 热数据不够从冷存储补 int need count - result.size(); ListLong 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 瓶颈或内存瓶颈再优化。过早抽象比不抽象更可怕因为你会花大量时间维护一套没人用的策略框架。希望帮到你。本文还有配套的精品资源点击获取
返回列表