“已读/未读”是职场IM(钉钉、飞书、企业微信)的标配功能。它让沟通变得确定——但也让“假装没看到”变得不可能。

单聊和群聊的已读回执,在技术难度上有着天壤之别。如果设计不当,一个5000人的大群,发一条消息就能让数据库瞬间写入5000条记录,并发ACK请求能把服务器打崩。
今天,我们就从基础实现开始,一步步演进到高并发架构,看看真正的工业级IM是如何扛住流量洪峰的。
| 特性 | 单聊(1v1) | 群聊(Group) |
|---|---|---|
| 送达回执 | 支持 | 不支持(群消息只存一份) |
| 消息已读回执 | 支持 | 支持(但实现极复杂) |
| 会话级批量已读 | 支持 | 支持(清空红点) |
| 存储模型 | 独立存储 | 一份消息,N份进度 |
单聊的已读逻辑非常直观:A发消息 → B读消息 → B告诉服务器“我读了” → 服务器通知A。
表结构设计(单聊消息表):
复制代码CREATE TABLE single_msgs (
msg_id BIGINT AUTO_INCREMENT PRIMARY KEY,
from_uid BIGINT NOT NULL,
to_uid BIGINT NOT NULL,
content TEXT,
send_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
is_read BOOLEAN DEFAULT FALSE, -- 是否已读
read_time TIMESTAMP NULL -- 已读时间
);
客户端流程:
ackConversationRead(to_uid) → 批量更新所有未读消息的 is_read = TRUE。onMessageRead(msg_ids),将UI状态从“已送达”刷新为“已读”。群聊的核心痛点是:群消息只有一份,但每个成员的阅读进度各不相同。
假设群里有200人,消息ID为10086,有的人读到了10090,有的人只读到了10050。
最朴素(但致命)的设计——每条消息为每个成员存一条回执明细:
复制代码-- ️ 低并发可用的基础表,高并发会崩溃!
CREATE TABLE msg_read_detail (
id BIGINT AUTO_INCREMENT PRIMARY KEY,
msg_id BIGINT NOT NULL, -- 群消息ID
user_id BIGINT NOT NULL, -- 群成员ID
is_read BOOLEAN DEFAULT FALSE,
read_time TIMESTAMP NULL,
UNIQUE KEY uk_msg_user (msg_id, user_id)
);
当发送方点击“查看已读详情”时:
复制代码SELECT COUNT(*) FROM msg_read_detail WHERE msg_id = 10086 AND is_read = TRUE;
SELECT user_id FROM msg_read_detail WHERE msg_id = 10086 AND is_read = FALSE;
这个方案在100人以内的小群里跑得挺好,但一到高并发环境,三天之内必崩。 为什么?接着往下看。
假设一个2000人的全员大群,早高峰时一条@所有人的消息发出:
若采用基础明细表,发1条消息 = 插入 1999条 msg_read_detail 记录。一天20条@消息 = 4万条记录。一个月120万,一年1500万。MySQL很快被打爆。
当2000人同时读完消息,瞬间产生2000个并发请求去更新同一条消息的“已读计数”或更新自己的进度。数据库行锁导致平均响应时间从10ms飙升到3秒,接口超时雪崩。
2000人×(上报ACK请求 + 服务端向发送方推送“张三已读”通知)。若发送方也在这个大群,消息扩散系数呈指数级增长,WebSocket/长连接带宽瞬间拉满。
核心思想:不再为每条消息存储每个成员的阅读记录,而是为每个用户在每个群只存一个值——「该用户在该群已读的最大消息ID(max_read_msg_id)」。
用户.max_read_msg_id >= 10086。存储介质选型:
| 存储层 | 职责 | 特点 |
|---|---|---|
| Redis(主) | 扛所有实时读写 | 内存操作,10万+ QPS |
| MySQL(从) | 冷备/历史归档 | 定时异步刷盘,5秒延迟可接受 |
Redis数据结构(Hash) :
复制代码# Key: group:read:{groupId}
# Field: {userId}, Value: {maxReadMsgId}
HSET group:read:123 456 10086 # 用户456在群123已读到10086
HGET group:read:123 456 # -> 10086
MySQL备份表(异步同步):
复制代码CREATE TABLE group_read_watermark (
group_id BIGINT NOT NULL,
user_id BIGINT NOT NULL,
max_read_msg_id BIGINT NOT NULL DEFAULT 0,
update_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (group_id, user_id) -- 每个用户在每个群只有一条!
);
效果:写操作从 O(N)(N=群人数)降为 O(1) 。存储成本骤降99.9%!
高并发下,网络延迟可能导致旧消息的ACK比新消息的ACK后到达。必须保证只向更大的MsgId更新:
复制代码-- Redis Lua脚本(原子执行)
local key = KEYS[1] -- group:read:{groupId}
local user_id = ARGV[1]
local new_msg_id = tonumber(ARGV[2])local current = redis.call('HGET', key, user_id)
if not current or new_msg_id > tonumber(current) then
redis.call('HSET', key, user_id, new_msg_id)
return 1 -- 更新成功
end
return 0 -- 忽略旧ID
所有写请求不直连DB,走Redis → Kafka/RocketMQ → MySQL:
(groupId, userId, msgId) 投递到Kafka。REPLACE INTO 或 INSERT ... ON DUPLICATE KEY UPDATE 写入MySQL。客户端聚合(Batching) :
客户端不逐条上报,而是维护本地水位线:
{ groupId, maxReadMsgId: 最新ID }。服务端限流(Throttling) :
对同一用户ID进行限频,例如 1秒内只允许上报2次,多余的请求直接返回200 OK(伪成功)并丢弃,因为最终一致性允许几秒钟的延迟。
若要展示“X人已读”,不要遍历Hash全量成员(大群会有性能灾难)。
正确姿势:在消息发出时预创建计数器。
SET read_count:{msgId} 0。INCR read_count:{msgId}。GET read_count:{msgId},O(1)时间复杂度。 复制代码// 伪代码:上报已读时原子增加计数
public void reportRead(String groupId, String userId, Long msgId) {
String watermarkKey = "group:read:" + groupId;
// 1. 更新水位线(Lua)
boolean updated = updateWatermark(watermarkKey, userId, msgId);
if (updated) {
// 2. 若更新成功,计数器+1(仅对近期消息有效,设置过期时间)
String countKey = "read_count:" + msgId;
redisTemplate.opsForValue().increment(countKey);
redisTemplate.expire(countKey, 7, TimeUnit.DAYS); // 7天后自动淘汰
}
}
不可能一套方案打天下,必须按群人数分治:
| 群规模 | 存储策略 | 已读列表展示 | ACK策略 |
|---|---|---|---|
| 小群 (2~100人) | 允许写明细表(成本可控) | 展示全部成员详细状态 | 实时 |
| 中群 (100~500人) | Redis水位线 + MySQL备份 | 仅显示“X人已读”,点击详情展示前20个最近活跃用户 | 客户端聚合5秒 |
| 大群 (500~3000人) | 纯Redis水位线 | 仅展示统计数字,不展示具体人头 | 强限流+10秒聚合 |
| 超大群/直播群 (>3000人) | 彻底关闭已读回执 | 不展示任何已读状态 | 不上报ACK |
高并发分布式系统下,异常是常态,必须优雅容错:
消息ID乱序(Old ACK覆盖New ACK) :
Max 比较机制解决,旧ID永远无法覆盖新ID。Redis突发故障(缓存雪崩/宕机) :
用户退群/重新加群:
group:read:{groupId} 中该用户的Field(释放内存)。max_read_msg_id 初始化为当前群最新消息ID。历史消息全部视为“已读”,不再追溯。(产品上合理,因为新成员没看过历史。)消息过期淘汰:
1. 服务端接收ACK上报(高并发入口)
复制代码@Service
@Slf4j
public class ReadReceiptService { @Autowired
private RedissonClient redissonClient;
@Autowired
private KafkaTemplate<String, String> kafkaTemplate; // 限流器(每用户每秒最多2次)
private final LoadingCache<String, RateLimiter> limiterCache = Caffeine.newBuilder()
.expireAfterWrite(1, TimeUnit.MINUTES)
.build(key -> RateLimiter.create(2.0)); public void reportRead(String groupId, String userId, Long msgId) {
// 1. 限流(防刷)
String limitKey = "ack:" + userId;
if (!limiterCache.get(limitKey).tryAcquire()) {
log.warn("User {} ack too frequently, dropped", userId);
return; // 静默丢弃
} // 2. 原子更新水位线(Lua)
String key = "group:read:" + groupId;
RScript script = redissonClient.getScript();
String lua = "local cur=redis.call('hget', KEYS[1], ARGV[1]);" +
"if not cur or tonumber(ARGV[2]) > tonumber(cur) then " +
"redis.call('hset', KEYS[1], ARGV[1], ARGV[2]); return 1; end; return 0;";
Long updated = script.eval(RScript.Mode.READ_WRITE, lua,
RScript.ReturnType.INTEGER,
Collections.singletonList(key),
userId, String.valueOf(msgId)); // 3. 更新计数器(仅当水位提升时)
if (updated != null && updated == 1L) {
String countKey = "read_count:" + msgId;
redisTemplate.opsForValue().increment(countKey);
redisTemplate.expire(countKey, 7, TimeUnit.DAYS);
// 4. 异步投递Binlog给Kafka(最终持久化到MySQL)
kafkaTemplate.send("topic_read_binlog",
groupId + "|" + userId + "|" + msgId);
}
}
}
2. 查询已读人数(高频读)
复制代码public Long getReadCount(Long msgId) {
String key = "read_count:" + msgId;
Integer count = (Integer) redisTemplate.opsForValue().get(key);
if (count != null) {
return Long.valueOf(count);
}
// 降级:查MySQL(但仅对近期消息)
return readCountMapper.selectCountByMsgId(msgId);
}
3. Kafka消费者(异步刷盘)
复制代码@Component
@Slf4j
public class ReadBinlogConsumer { @Autowired
private JdbcTemplate jdbcTemplate; @KafkaListener(topics = "topic_read_binlog", batch = "true")
public void consume(List<String> records) {
// 批量去重:同一个(groupId,userId)只取最大的msgId
Map<String, Long> latestMap = new HashMap<>();
for (String record : records) {
String[] parts = record.split("|");
String key = parts[0] + ":" + parts[1]; // groupId:userId
Long msgId = Long.parseLong(parts[2]);
latestMap.merge(key, msgId, Math::max);
} // 批量Replace Into MySQL (减少锁竞争)
String sql = "INSERT INTO group_read_watermark (group_id, user_id, max_read_msg_id) " +
"VALUES (?, ?, ?) ON DUPLICATE KEY UPDATE max_read_msg_id = ?";
List<Object[]> batchArgs = new ArrayList<>();
for (Map.Entry<String, Long> entry : latestMap.entrySet()) {
String[] ids = entry.getKey().split(":");
batchArgs.add(new Object[]{
Long.parseLong(ids[0]),
Long.parseLong(ids[1]),
entry.getValue(),
entry.getValue()
});
}
jdbcTemplate.batchUpdate(sql, batchArgs);
log.info("Flushed {} read states to MySQL", batchArgs.size());
}
}
| 维度 | 低并发方案(基础) | 高并发方案(进阶) |
|---|---|---|
| 存储模型 | (msg_id, user_id) 明细 | (group_id, user_id) 水位线 |
| 存储介质 | MySQL单库 | Redis(主) + MySQL(从) |
| 写操作复杂度 | O(N) | O(1) |
| ACK上报方式 | 逐条实时上报 | 聚合批量(10条/5秒) |
| 已读计数查询 | 全表扫 COUNT | 原子计数器 INCR/GET |
| 最终一致性 | 强一致(实时) | 秒级最终一致 |
| 大群支持 | ≤100人 | 万人级(需分级降级) |
希望这篇从基础到进阶的硬核实战能帮你避开那些“血泪坑”。如果在落地中遇到更极端的场景(如跨境网络延迟、多机房同步),欢迎评论区一起探讨!