
简介这份压缩包是一套基于Java的Timeline模式抽象库面向社交与IM类应用开发者覆盖朋友圈/微博式时间线、消息推送、feed流以及基础IM通讯核心价值在于提供数据流间的统一分发能力帮助降低后端动态信息实时同步与排序检索的开发门槛。包内共71个文件主体为51个Java源码文件搭配9个XML配置、2个FreeMarker模板、1个YAML、SQLite数据库文件及说明文档整体约1.59MB。项目按timeline-core、timeline-store、timeline-example、timeline-starter等模块组织其中core实现核心分发逻辑store提供Redis与内存两种存储实现starter便于Spring Boot接入example内含IoT和IM两种示例场景可直接参照搭建消息分发链路或扩展自己的业务模块。目前已有84人学习下载适合具备一定Java基础、希望快速落地社交feed或消息推送能力的开发者参考。1. 把朋友圈、微博、IM 收进同一套 Java 抽象层Timeline 模式到底在抽象什么做过社交或消息类业务的人基本都撞过同一堵墙朋友圈时间线、微博关注流、IM 聊天记录、系统消息推送四套业务长得不一样底层逻辑却惊人一致——都是“某个用户产生内容分发给一批关注者按时间倒序呈现”。我在接手一个社区项目时发现代码里有三套各自为政的查询逻辑Feeds 服务查 MySQL、消息中心查 Redis 列表、IM 模块干脆自己维护一套内存队列维护成本高到让人想跑路。后来把数据模型统一收敛到 Timeline 抽象层才真正解决这个问题。这份基于 Java 实现的抽象库核心就是把 feed 流、朋友圈、IM 通讯、消息推送抽象成同一种“数据流 分发”模型生产者产生事件抽象层按关注关系写扩散或读扩散到目标用户的时间线里消费者按游标拉取增量。适用面很广社区信息流、即时通讯未读消息、运营推送、好友动态聚合都能直接拿这套模型落地。适合正在设计 feed 流存储、或者想把多个消息型业务收敛到统一架构的 Java 从业者。下面我从存储模型、分发设计、IM 接入和踩坑几个角度逐层拆代码都是可直接搬的骨架。2. Timeline 的存储模型收件箱、发件箱与推拉结合的选择2.1 为什么收件箱 / 发件箱分离是 feed 流的基石Timeline 模式的核心不是“时间”本身而是数据组织方式。你在朋友圈看到的好友动态并不是每次请求都现场去好友表里查一遍而是提前写到“属于你的收件箱”里。所谓收件箱inbox是每个用户维度下的一条按时间排序的 feed 列表发件箱outbox则是当前用户自己发布的内容列表。关注关系变化时发布的内容要同步到所有粉丝的 inbox这里就产生了写放大和读放大的取舍。写扩散fanout-on-write是指用户发一条动态直接把内容推送到所有粉丝的收件箱读取时只查自己的 inbox速度快、延迟低适合微博、朋友圈这种关注数有限、读多写少的场景。读扩散fanout-on-read则不推粉丝刷新时动态去关注列表聚合适合大 V 账号——粉丝上百万时写扩散直接把存储打爆。这份 Java 抽象库的做法是支持配置切换默认走写扩散同时对超大粉丝量账号走读扩散混合这也是目前业界的主流折中方案。这段取舍直接决定抽象层的接口设计。我在库内部用 TimelineStorage 接口把存储细节隔离掉对外只暴露 append、remove、fanout 三个方法MySQL、Redis、本地内存各实现一套业务侧完全不感知底层是什么。下面这个代码块是抽象库的核心存储接口骨架public interface TimelineStorage { // 向用户的某个时间线追加一条消息/动态 boolean append(String userId, String timelineType, TimelineEntry entry); // 从时间线移除一条消息如删除动态、撤回IM消息 boolean remove(String userId, String timelineType, String entryId); // 分页拉取时间线游标用 lastId 而不是 offset PageResultTimelineEntry pull(String userId, String timelineType, String lastId, int limit); // 批量写入用于写扩散的扇出场景 void fanout(CollectionString targetUsers, String timelineType, TimelineEntry entry); }append 负责写入时间线timelineType 区分 timeline:feed、timeline:im、timeline:push 等业务域pull 用 lastId 做游标分页而不是页码因为 feed 流在刷新过程中不断有新数据进来用 offset 会导致重复或漏数据fanout 是写扩散入口内部会遍历粉丝列表逐个写入。库默认使用 Redis 的 ZSet 实现 TimelineStoragescore 用毫秒时间戳member 用唯一的消息 ID这样拉取时按 score 倒序取一批即可删除时直接 ZREM。如果你要接 MySQL就把 timelineType 加进分表键按 user_id 哈希分 64 张表。2.2 分发器的动作拆分扇出、合并与削峰抽象库的第二个关键类是 Dispatcher负责把一条内容变成一次扇出动作。实际场景里不能简单 for 循环写粉丝列表一是粉丝数量大循环写会造成数据库压力尖峰二是要支持“部分粉丝不感兴趣”“自己发的动态只进自己时间线”“评论不扇出”这类业务过滤三是写扩散失败要能补偿。所以 Dispatcher 内部做三步先将一条消息解析成 FanoutTaskFanoutTask 携带目标用户列表、负载数据和渠道类型再通过任务队列库内支持线程池或 Redis Stream分批消化最后对写失败的待办用户做补偿记录到一个 retry topic由定时任务重放。这部分的代码抽象和配置参数很值得直接抄我给出 Dispatcher 的核心实现public class Dispatcher { private final TimelineStorage storage; private final ExecutorService fanoutPool; private final int batchSize; public Dispatcher(TimelineStorage storage, int fanoutThreads, int batchSize) { this.storage storage; this.fanoutPool new ThreadPoolExecutor(fanoutThreads, fanoutThreads, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue(10000)); this.batchSize batchSize; } // 传入原始事件如发布动态、发送IM消息由Dispatcher决定扇出策略 public void dispatch(FeedEvent event, SetString followers) { ListString fanoutList filterSpecialUsers(followers); // 按batchSize切分避免一次性创建太多任务 ListListString partitions Lists.partition(fanoutList, batchSize); for (ListString partition : partitions) { fanoutPool.submit(() - { for (String userId : partition) { storage.append(userId, event.getTimelineType(), event.toEntry()); } }); } } }这里有两个设计细节值得划重点。filterSpecialUsers 是业务扩展点——你可以在这里屏蔽黑名单、跳过自己、过滤不感兴趣的用户batchSize 参数控制批大小默认 500如果粉丝超过十万务必配合 MQ 做异步削峰而不是在请求线程里直接扇出否则发布一条动态会让接口耗时从 10ms 涨到 300ms。fanoutThreads 控制并发写线程建议设置为磁盘 IO 能力的 2 到 3 倍不要盲目调大否则 Redis 连接池先被打满。2.3 拉取路径的抽象时间线聚合、未读过滤与分页分发写好了读取侧同样要抽象。用户的首页时间线是多个来源的聚合好友动态、自己发的、系统推送的、IM 消息提醒不同来源的权重也不同。抽象库的做法是允许业务注入一个 TimelineAggregator它负责解释 TimelineQuery 参数最终调用 TimelineStorage.pull。注意这里的读取不是简单的一层获取而是按时间线类型合并拉取再做一次内存归并排序。public class TimelineQuery { private String userId; private ListString timelineTypes; // 要拉取的渠道列表 private long lastTimestamp; // 游标时间精确到毫秒 private int maxCount; // 每次最大拉取条数默认20 private boolean filterRead; // 是否过滤已读IM场景常用 private BiFunctionTimelineEntry, String, Boolean visibleFilter; // 业务可见性过滤 // 省略getters/setters } // 聚合逻辑示例多种时间线归并后按时间倒序截断 public ListTimelineEntry mergePull(TimelineQuery query) { ListPageResultTimelineEntry pages new ArrayList(); for (String type : query.getTimelineTypes()) { PageResultTimelineEntry page storage.pull(query.getUserId(), type, String.valueOf(query.getLastTimestamp()), query.getMaxCount() * 2); pages.add(page); } return Stream.concat(pages.stream().flatMap(p - p.getList().stream())) .sorted(Comparator.comparingLong(TimelineEntry::getTimestamp).reversed()) .limit(query.getMaxCount()) .collect(Collectors.toList()); }mergePull 里用 2 倍数量预取是个小技巧——因为多源归并后要截断到 maxCount如果只取 maxCount 条排序后可能某个渠道一条都进不了首页导致内容单一按 2 倍预取后截断能让各渠道内容相对均匀。maxCount 建议设置上限 50这是移动端的合理首屏数量。filterRead 用于 IM 未读场景由读取侧把收到的消息标记为已读抽象库不自动处理因为已读状态属于业务语义。3. 数据流分发机制订阅-发布模型与广播的 Java 实现3.1 事件总线与数据流分发的拓扑结构数据流之间的分发是抽象库区别于普通时间线组件的关键能力。简单说一条朋友圈动态不仅要落到发布者自己的发件箱还要触发三方动作推送通知到 APNs/极光、更新粉丝收件箱、发送 WebSocket 消息给在线的铁粉。这就是多分发拓扑。抽象库把这个能力设计成事件总线模式所有业务动作统一封装成 TimelineEvent总线根据事件类型路由到不同的 Dispatcher 和 Listener。传统做法是代码里到处调用 service发完动态调 pushService.push()、imService.send()、feedService.update()调用关系混乱且试错成本高。用事件总线后发布者只发布一个 Event路由完全解耦。Java 里实现这个总线不需要引入 Spring 事件那套自己用 ConcurrentHashMap 维护 topic 与订阅者列表即可。核心分发接口长这样public interface EventListenerT extends TimelineEvent { // 返回感兴趣的事件类型如 feed.publish、im.message String subscribeTopic(); void onEvent(T event); } public class TimelineEventBus { private final ConcurrentMapString, ListEventListener? subscribers new ConcurrentHashMap(); public T extends TimelineEvent void register(EventListenerT listener) { subscribers.computeIfAbsent(listener.subscribeTopic(), k - new CopyOnWriteArrayList()) .add(listener); } SuppressWarnings(unchecked) public T extends TimelineEvent void publish(T event) { ListEventListener? listeners subscribers.get(event.getTopic()); if (listeners null) return; for (EventListener? listener : listeners) { try { ((EventListenerT) listener).onEvent(event); } catch (Exception e) { // 单条监听器异常不能阻断其他分发必须捕获 log.error(event dispatch failed, topic{}, event.getTopic(), e); } } } }同步分发的好处是代码简单、调试直观缺点是慢。如果 onEvent 里要做网络调用整个发布链路会被拖慢。所以我的习惯用法是TimelineEventBus 负责同步结构分发但每个 listener 内部自己决定是异步还是同步。比如 IM 推送 listener 里丢线程池状态更新 listener 保持同步。subscribeTopic 用字符串层级命名比如 feed.publish、im.message、system.notice由业务方规定总线不关心具体语义这样就保持了抽象库的通用性。3.2 分发失败的重试与死信处理数据流分发最容易翻车的不是主流程而是重试逻辑。我最初用的方案是失败就 catch 住打日志表面看系统很稳定实际丢了一堆推送没人察觉。后来在抽象库里加了 RetryTask 机制分发失败的 listener 可以选择抛出 RetryableException总线会把事件封装成 RetryTask写入一个内部延迟队列按指数退避重试 3 次超过次数进入死信列表由运维脚本人工处理。public class RetryTask { private final String taskId; private final TimelineEvent event; private final int retryCount; private final long nextExecuteTime; public RetryTask(TimelineEvent event, int retryCount) { this.event event; this.retryCount retryCount; this.taskId UUID.randomUUID().toString(); // 指数退避30s - 60s - 120s this.nextExecuteTime System.currentTimeMillis() 30000L * (1 retryCount); } public boolean needRetry() { return retryCount 3; } }这个类的关键是 nextExecuteTime 的计算1 retryCount 做指数退避初学者很容易写成固定间隔遇到下游抖动会直接雪崩。注意死信不是丢了不管而是要落到 Redis 的 zset 里按执行时间排序方便手动查询和重放。还要注意重试的幂等同一个事件重放时不能造成重复插入方案是在 TimelineEntry 里带上 eventId写入时用 Redis 的 SETNX 做幂等控制或者让数据库表对 (user_id, timeline_type, event_id) 建唯一索引。3.3 在线状态感知IM 消息走实时通道feed 流走拉取通道分发拓扑里IM 消息和 feed 流的实时性要求完全不同。IM 消息必须毫秒级到达feed 流秒级可接受。所以抽象库把分发通道分成两条在线通道用 WebSocket 推离线通道写收件箱等待拉取。这需要分发器感知用户在线状态库内抽象一个 PresenceService 接口业务方可接自己的长连接网关。在线时消息既写收件箱保证多端同步又通过 WebSocket 实时下发离线时只写收件箱待用户上线后客户端拉取未读。public class ImMessageDispatcher extends AbstractDispatcher { private final WsGateway wsGateway; private final PresenceService presenceService; Override public void dispatch(ImMessage message, SetString receivers) { for (String receiver : receivers) { // 收件箱写入是必须的保证离线不丢消息 storage.append(receiver, timeline:im, message.toEntry()); if (presenceService.isOnline(receiver)) { wsGateway.send(receiver, message.toPushPayload()); } else { // 离线走推送服务如APNs、极光、厂商通道 pushService.push(receiver, message.toPushPayload()); } } } }这里有个容易被忽略点写收件箱和 WS 推送不是原子操作先写库再推还是先推再写库都有窗口。我的经验是先写库后推送因为推送失败最多是用户多点一次刷新如果先推后写库用户可能收到消息但拉取历史记录时找不到体验更差。如果对一致性要求极高需要把两个动作放进同一个本地消息事务表由另外的 worker 保证最终一致但大多数业务场景没必要付出这个成本。4. 接入实战把抽象库用在一个模拟朋友圈 IM 通讯项目中4.1 场景定义与依赖引入方式抽象库本身高度耦合业务直接落代码前必须先定场景。假设我们要做一个简化版社交 App包含三个业务域好友发布动态后粉丝可见、单聊消息实时收发、系统公告推送。对应到抽象库就是三种 Timeline 类型和三条分发链路。引入方式上库以 jar 形式提供核心依赖只有一个 slf4j-api存储层由业务侧注入 RedisTemplate 或 JDBC 实现这样设计是为了不绑架使用方。下面是典型的 Maven 坐标和初始化代码dependency groupIdcom.example/groupId artifactIdtimeline-abstract-lib/artifactId version2.1.0/version /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency初始化时需要把存储实现、分发线程池和事件总线装配起来。注意这里的版本号是我项目里用的你使用时应替换为实际下载包里的版本。装配代码如下Configuration public class TimelineConfig { Bean public TimelineStorage timelineStorage(StringRedisTemplate redisTemplate) { return new RedisTimelineStorage(redisTemplate); } Bean public TimelineEventBus eventBus() { return new TimelineEventBus(); } Bean public Dispatcher dispatcher(TimelineStorage storage) { // 线程数根据压测调整不要盲目大 return new Dispatcher(storage, 8, 500); } }装配阶段最容易翻车的点是 RedisTemplate 的序列化器。默认 JDK 序列化在 key 里会产生 \xAC\xED 前缀导致 ZSet 操作时 key 匹配不上我第一周就吃过这个亏。务必把 key 设置为 StringRedisSerializervalue 用 Jackson 或 Fastjson尤其是 entryId 直接用 String 类型存储。4.2 发布一条朋友圈的完整调用链从发布到粉丝可见完整链路是客户端 POST /feed/publish → 业务层创建 FeedEvent → 事件总线发布 feed.publish → FeedPublishListener 将动态写入发布者发件箱 → 同时获取关注者列表 → Dispatcher 扇出到所有粉丝收件箱 → 异步推送通知给在线粉丝。这个链路里发布者的接口只感知到事件总线 publish 这一步后续全是异步接口 RT 可以稳定在 20ms 内。代码实现如下public class FeedService { private final TimelineEventBus eventBus; private final TimelineStorage storage; public String publish(FeedPublishRequest request, String userId) { String feedId IdGenerator.next(); // 先写自己发件箱保证发布立即可见 TimelineEntry ownEntry TimelineEntry.builder() .entryId(feedId) .content(request.getContent()) .timestamp(System.currentTimeMillis()) .build(); storage.append(userId, timeline:own, ownEntry); // 发布事件粉丝可见性由监听器负责 FeedEvent event FeedEvent.builder() .topic(feed.publish) .feedId(feedId) .authorId(userId) .content(request.getContent()) .build(); eventBus.publish(event); return feedId; } } // 监听器内部做扇出 Component public class FeedPublishListener implements EventListenerFeedEvent { private final Dispatcher dispatcher; private final FollowService followService; Override public String subscribeTopic() { return feed.publish; } Override public void onEvent(FeedEvent event) { SetString followers followService.getFollowers(event.getAuthorId()); dispatcher.dispatch(event, followers); } }为什么发件箱要同步写、粉丝收件箱要异步写因为发布者自己刷新时必须立刻看到动态如果也异步会出现“发布成功后刷新看不到”的尴尬而粉丝侧延迟几秒完全无感知异步换来了发布接口的稳定性。如果你的业务要求粉丝毫秒级可见可以调整 Dispatcher 的 fanoutPool 线程数和 batchSize1000 粉丝以内写扩散基本 100ms 内完成。4.3 IM 消息的分发与未读计数IM 场景和 feed 的差异在于IM 消息需要未读红点、已读回执和时序唯一。抽象库对 IM 消息的处理是消息 ID 用发件人 毫秒时间戳 自增序号组合生成确保同一会话内严格递增消息写入收件箱时 score 即是消息序号保证拉取顺序稳定。未读数通过 zset 的 score 游标实现——客户端记录 lastReadSeq拉取时所有 score 大于 lastReadSeq 的 member 计数就是未读数不需要单独维护 counter 字段这一招能省掉很多一致性问题。public long getUnreadCount(String userId, String peerUserId) { String timelineKey timeline:im: userId : peerUserId; String lastReadSeqStr stringRedisTemplate.opsForValue().get(lastread: userId : peerUserId); long lastReadSeq lastReadSeqStr null ? 0L : Long.parseLong(lastReadSeqStr); // ZSet中score即消息序号返回大于lastReadSeq的数量 return stringRedisTemplate.opsForZSet().count(timelineKey, lastReadSeq 1, Long.MAX_VALUE); }单聊会话的时间线 key 是 timeline:im:{userId}:{peerUserId}这样双方各自有一条独立的会话时间线语义清晰删除会话时直接 DEL 该 key不会影响全局未读计算。ZSet count 的时间复杂度是 O(log N)无论会话多长未读数查询都不受影响。注意这里的 lastReadSeq 要单独存不要存在 TimelineEntry 里否则每次拉消息都要扫描整条时间线。5. 避坑排查我从这套抽象库落地中总结的六个真实翻车现场5.1 时间线分页出现重复数据现象客户端上拉刷新时首页时不时出现一条已经看过的动态位置还总在顶部附近。原因一开始我图省事用 offset 分页feed 流在消费者刷新过程中不断有新数据写入offset 会整体后移导致上一页尾部数据被挤进下一页造成重复。这也是抽象库把 push 接口设计成 lastId 游标的原因。解决完全放弃 offset所有分页参数统一改为 lastId即上一页最后一条动态的 score毫秒时间戳。第一次刷新 lastId 传 0后续传接口返回的 lastTime 字段。注意 lastId 不能只传时间戳还要带上毫秒内的自增序列否则同一毫秒产生的多条动态会漏。5.2 fanout 大批量写入把 Redis 连接池打满现象某大 V 发布一条动态后应用日志里全是 Redis connection timeout其他业务接口也跟着变慢持续数分钟才恢复。原因Dispatcher 线程池开到 32batchSize 又设置成 2000大 V 有 200 万粉丝等于瞬间提交了 1000 个扇出任务Redis 连接池上限只有 50所有线程都在等待获取连接。解决三层控制——fanout 线程数收敛到 8batchSize 降到 500给 fanout 写入操作单独开一个 Redis 连接池与业务读写隔离。还有一个血泪经验对粉丝数超过 10 万的账号直接走读扩散不写粉丝收件箱改为粉丝拉取时去聚合该大 V 的发件箱存储量能省 90%。5.3 事件总线监听器异常导致主链路中断现象发布动态偶尔变成 500 错误查日志发现是推送监听器里调第三方推送接口超时把整个 publish 接口拖垮。原因TimelineEventBus 的 publish 是同步遍历 listener并且最开始没有 catch 异常一个 listener 抛错直接中断后续所有 listener 执行。解决两个改动——总线遍历 listener 时务必 try/catch单点异常仅记录日志保证后续分发不中断需要强一致性的核心监听器如写收件箱放到总线调用链最前面非核心的推送监听器用 Async 丢到独立线程池彻底物理隔离。5.4 IM 消息撤回后未读数不正确现象用户撤回一条已读消息另一端的未读红点数字没变或者 A 发两条消息撤回一条B 的未读数变成 -1。原因撤回操作直接调 ZREM 从时间线里删除了 member但 ZSet 的 count 统计是按 score 区间的删除后 lastReadSeq 没更新统计的区间范围仍然包含已删除的序号导致计数错乱。解决撤回消息不要物理删除而是在 TimelineEntry 里加 revoked 标记字段拉取时过滤掉未读数统计改为基于消息序号而不是实际 member 数量。正确做法是维护一个 sessionSeq 自增数每次发消息前先 INCR撤回时记录 revokedSeq 集合未读数 sessionSeq - lastReadSeq - revokedCountBetween(lastReadSeq, sessionSeq)。5.5 多端登录时消息重复消费现象用户同时在手机和 PC 登录WS 通道把同一条 IM 消息推送两次客户端展示出现重复气泡。原因抽象库的在线通道分发只判断用户是否在线没区分端。手机和 PC 各一条连接都满足 isOnline 条件都收到消息。解决wsGateway.send 里增加 deviceId 维度同一用户多条连接时只推送给最后活跃的连接设备其他设备通过拉取收件箱同步。注意这里要区分“多端在线”和“多端同步”前者的设计目标是实时推一端的通知后者的数据一致性由收件箱拉取兜底。5.6 本地测试时始终拉不到新数据现象本地起服务调试发布动态后 Redis 里能查到数据但接口返回的列表始终是旧的。原因RedisTimelineStorage 的 pull 方法用了缓存提高查询性能——先查本地缓存没命中才查 Redis缓存过期时间设了 60 秒导致新数据一直被缓存挡住。解决开发环境把缓存过期时间设为 0也就是完全旁路缓存生产环境对首页这种实时性要求高的场景缓存时间也别超过 5 秒或者干脆只在热点 Key 上开缓存。抽象库的 Storage 接口加了一个 开启缓存开关 配置我在本地调试时默认关闭上线前再按规则打开。6. 压测验证这套抽象库从单机到集群的三种验证姿势资源拆解完最重要的还是验证。我推荐三条递进路径单元验证 → 单机压测 → 集群长稳。单测阶段只验证语义不需要真实 Redis用本地内存版 TimelineStorage 跑全链路重点抓业务逻辑正确性这一步可以在不依赖环境的情况下完成大部分 debug。内存版实现里要注意 ConcurrentHashMap 的并发写安全模拟多线程扇出时用 CopyOnWriteArrayList 或者加锁否则会有 ConcurrentModificationException。接下来是单机压测。我见过太多人直接上集群压测结果环境因素干扰判断连数据结构合理性都验证不了。正确做法是先在没有网络抖动干扰的情况下固定一台机器压发布接口与拉取接口观察耗时分布。压测时注意一个关键参数每次拉取 20 条内数据响应时间应该稳定在 20ms 内才说明存储层没问题如果波动大大概率是 Redis 连接池配置过小或者 fanout 线程池队列溢出。执行压测的脚本我习惯用 Gatling 或者 Jmeter但更直接的方式是写一段 Java 刚性压测代码循环调 Dispatcher 和 pull不做任何额外封装方便在 IDE 里断点调试public class SmokeLoadTest { public static void main(String[] args) throws Exception { TimelineStorage storage new MemoryTimelineStorage(); Dispatcher dispatcher new Dispatcher(storage, 8, 100); // 模拟一个用户产生动态并扇出给5000个粉丝 long start System.currentTimeMillis(); for (int i 0; i 5000; i) { FeedEvent event FeedEvent.builder() .topic(feed.publish) .feedId(feed- i) .authorId(user-1) .content(content- i) .build(); dispatcher.dispatch(event, Set.of(fan- i)); } System.out.println(5000条扇出完成耗时: (System.currentTimeMillis() - start) ms); // 验证时间线拉取结果 PageResultTimelineEntry page storage.pull(fan-4999, timeline:feed, 0, 20); System.out.println(粉丝4999拉取条数: page.getList().size()); } }这段压测代码的关键在于启动时人为控制变量——单机单存储、固定粉丝数、固定内容长短得到的数据才有比较意义。5000 条扇出的耗时在你的电脑上跟我这里可能不一样但如果超过 3 秒优先排查线程池阻塞和序列化耗时而不是拉代码级别的问题。集群长稳验证时重点观察三个指标事件总线消费积压数、Redis 内存增长曲线、拉取接口 P99 延迟。积压数持续上升说明消费能力不足先加 Dispatcher 线程再加分片内存增长过快则要检查是否有不合理的 fanout 数据残留和未设置过期时间的时间线IM 单聊的 ZSet 建议 30 天定期清理一次。最后加一个真实链路演练模拟用户下线再上线的场景验证离线消息从收件箱拉取是否完整这是我每次上线前强制走一遍的流程现在已经成为习惯——先把内存版跑通再压真实 Redis最后完整操作一遍全流程确认无异常才进入发布流程。从那次被推送超时坑过一次之后我所有涉及分发的项目都强制走这条验证路径希望帮到你。本文还有配套的精品资源点击获取