Java实现未读消息案例:从Redis到WebSocket的全链路架构解析
目录导读
- 未读消息系统的核心挑战与设计原则
- 基于Redis的未读计数存储方案
- 消息推送与实时更新:WebSocket + STOMP实践
- 消息已读/全部已读的幂等处理
- 数据库与缓存的一致性保障
- 高并发场景下的性能优化(批量标记、异步落库)
- 经典问答:消息丢失、重复推送、扩展性如何解决?
未读消息系统的核心挑战与设计原则
在IM、社交App或工单系统中,“未读消息”是高频核心功能,纯数据库COUNT查询在千万级数据量下会拖垮DB,且无法实时推送。核心设计原则是:读多写少用缓存,实时交互走长连接,最终一致性靠异步任务,典型的架构分层为:接入层(WebSocket) -> 业务层(Spring Boot) -> 存储层(Redis + MySQL)。

基于Redis的未读计数存储方案
Redis Hash结构是未读计数的首选,Key设计为 unread:{userId},Field为 conversationId(单聊)或 groupId(群聊),Value为未读数,额外设置一个全局总数字段 total,便于顶部红点展示。
代码示例(Spring Data Redis):
public void incrUnread(Long userId, Long conversationId) {
String key = "unread:" + userId;
redisTemplate.opsForHash().increment(key, conversationId.toString(), 1);
redisTemplate.opsForHash().increment(key, "total", 1);
}
关键优化:对群聊采用“批量写,延迟读”,用户不在线时,消息发到Redis的 pending:{userId} 列表,用户上线时一次性拉取并合并计数。
消息推送与实时更新:WebSocket + STOMP实践
当消息产生时,不仅要更新计数,还需实时推送给在线用户,使用Spring Boot集成WebSocket + STOMP,定义 /user/queue/unread 作为点对点通道。
推送逻辑:在 @MessageMapping 或 Service 层,通过 SimpMessagingTemplate.convertAndSendToUser(userId, "/queue/unread", unreadDto) 推送,前端收到后直接刷新角标数字,无需重新拉取全量列表。
防抖策略:如果短时间内消息暴涨(如群发),使用 @Scheduled 定时任务每500ms批量推送一次汇总计数,避免网络风暴。
消息已读/全部已读的幂等处理
用户点击某个会话时,调用 markAsRead(userId, conversationId),这里要处理并发点击与重复请求。
幂等实现:Redis的 SETNX 加锁,或者使用Lua脚本原子操作:
local unread = redis.call('HGET', KEYS[1], ARGV[1])
if unread and unread > 0 then
redis.call('HSET', KEYS[1], ARGV[1], 0)
redis.call('HINCRBY', KEYS[1], 'total', -unread)
end
return unread
需将“已读事件”异步写入MQ,由消费者更新MySQL会话表的 last_read_time,用于跨端同步。
数据库与缓存的一致性保障
Redis作为主读,MySQL作为持久层。落库时机:用户主动已读时,或定时任务每5分钟将 dirty 标记的未读变更刷入DB。
一致性方案:采用 Cache-Aside + 版本号,每次更新Redis时,自增 version;定时任务对比DB版本号,若不同则更新DB,防止缓存击穿与脏数据。
高并发场景下的性能优化
- 批量标记已读:当用户进入“全部已读”时,不要遍历所有conversationId,而是执行
SCAN获取全部Field,删除Key并重建total=0。 - 异步化:消息发送后,将写库操作丢入
ThreadPoolTaskExecutor,核心业务只操作Redis。 - 热点Key拆分:对于超大群(10万+人),用
unread:{groupId}:{userId%10}分片,防止单Key内存过大。
经典问答:消息丢失、重复推送、扩展性如何解决?
Q1: WebSocket断线重连后,未读计数会丢失吗?
A: 不会,WebSocket断开后,客户端会轮询REST接口 GET /unread/summary,该接口直接查Redis的Hash,重连成功后,Redis会再次推送全量计数,核心是Redis的AOF持久化+主从同步。
Q2: 用户在两台设备上操作,已读状态如何同步?
A: 采用 已读时间戳 而非布尔值,设备A标记已读时,写入 lastReadTime,设备B拉取时,对比该时间戳与消息时间,过滤出未读,Redis的 ZSET(score=时间戳)可以优雅解决。
Q3: 未读数一直不准,偶尔偏大或偏小,怎么排查?
A: 先查Redis unread:{userId} 的实际值,再核对MySQL的 unread_log 表,常见原因:① 未处理“消息撤回”导致的计数未减;② 重复消费MQ消息;③ 并发重置时未加事务锁。解决方案:在所有写操作上增加 operation_id 唯一键,消费端做去重表。
Java实现未读消息绝非简单的 UPDATE count = count + 1。核心是分层存储(Redis加速、MySQL兜底)、异步解耦(MQ削峰)、实时感知(WebSocket),对于中小型项目,上述架构足够支撑万级并发;对于大型IM系统,还需引入caffeine本地缓存和redisson分布式锁,建议开发者从“单一会话未读”起步,逐步扩展为“多维度红点系统(@我、系统通知、群@)”的通用组件,记住一句话:把计数放内存,把确认放持久层,把推送放长连接,把一致性放异步任务。