直播礼物的后端架构:消息推送、去重、幂等与分布式锁
直播礼物消息如何稳定、准确地送达直播间?本文拆解礼物消息推送、幂等去重、连击合批与分布式锁等核心机制,帮助开发者避免重复扣费、重复播放及高并发数据冲突。

本篇系统拆解直播礼物的后端架构 :礼物消息如何推送到直播间、如何去重防止重复扣费、如何保证幂等避免重复播放、以及分布式锁在礼物连击合批中的应用。
一、礼物消息推送架构

用户 A 在直播间送礼,主播和其他观众都要实时看到。核心诉求是低延迟 + 高并发。
典型架构:长连接网关 + 消息队列 + 推送服务
[客户端] --WebSocket--> [网关] --MQ--> [推送服务] --fanout--> [直播间所有客户端]长连接网关:每个客户端建立一条 WebSocket 连接,网关维护 userId -> connection 映射。网关只负责连接保活和消息转发,不处理业务逻辑。
消息队列:用户送礼后,业务服务把礼物消息写入 Kafka / Pulsar,推送服务订阅 MQ 消费。MQ 削峰、解耦,避免直接调用推送服务导致雪崩。
推送服务:从 MQ 消费礼物消息,查询直播间在线用户列表(Redis sorted set,score = 心跳时间戳),批量下发到网关。
消息格式
————————————————
{
"msgId": "gift_1234567890_001",
"roomId": "room_9527",
"senderId": "user_A",
"giftId": "gift_aurora_town",
"giftCount": 1,
"timestamp": 1735689600000,
"seq": 12345
}msgId 全局唯一,用于去重;seq 单调递增,用于客户端乱序检测和重排。
二、去重:防止重复扣费
送礼涉及扣费,重复扣费是红线。常见触发场景:
用户双击“送礼”按钮,客户端发两次请求;
网络抖动,客户端超时重试;
业务侧重复消费 MQ 消息。
方案一:客户端生成 requestId + Redis 去重
客户端每次点击送礼生成唯一 requestId(UUID),请求带上它。服务端用 Redis SET NX 写入:
————————————————
String key = "gift:dedup:" + requestId;
Boolean success = redis.setIfAbsent(key, "1", 10, TimeUnit.SECONDS);
if (!success) {
return Result.error("重复请求");
}
// 继续扣费逻辑10 秒过期,足够覆盖正常请求的生命周期。
方案二:数据库 唯一索引
礼物订单表加唯一索引 (userId, requestId),插入时若冲突则返回“已处理”:
INSERT INTO gift_order (user_id, request_id, gift_id, amount, created_at)
VALUES (?, ?, ?, ?, NOW())
ON DUPLICATE KEY UPDATE updated_at = NOW();数据库兜底,即使 Redis 失效也不会重复扣费。
三、幂等:防止重复播放
扣费去重之后,推送侧也要保证幂等——同一条礼物消息不能在客户端播两次。
客户端幂等检查
客户端维护 Set<String> playedMsgIds,收到消息时先判断:
if (playedMsgIds.contains(msg.msgId)) {
return; // 已播放,跳过
}
playedMsgIds.add(msg.msgId);
playQueue.offer(msg);playedMsgIds 设置大小上限(如 1000 条),超过后移除最老的(LRU)。
服务端去重表
推送服务在下发前写 Redis 去重表:
String key = "gift:push:" + roomId + ":" + msgId;
Boolean pushed = redis.setIfAbsent(key, "1", 60, TimeUnit.SECONDS);
if (!pushed) {
return; // 已推送,跳过
}
// 继续推送60 秒过期,覆盖直播间消息的典型生命周期。
四、分布式锁:礼物连击合批

用户在 3 秒内连续点击 10 次送礼,如果每次都单独推送一条消息,直播间会被刷屏。更好的方案是合批:3 秒窗口内的多次送礼合并成一条 giftCount=10 的消息。
场景:多台推送服务器同时处理同一用户的连击
用户 A 在 1 秒内点了 5 次,这 5 条消息可能被 MQ 分到不同的消费者(推送服务实例)。如果不加锁,会推 5 条独立消息到客户端。
分布式锁实现
String lockKey = "gift:batch:" + roomId + ":" + senderId + ":" + giftId;
RLock lock = redisson.getLock(lockKey);
try {
// 尝试加锁,最多等 100ms,锁持有时间 3s
if (lock.tryLock(100, 3000, TimeUnit.MILLISECONDS)) {
// 读取当前窗口内已累计的 count
String countKey = "gift:batch:count:" + roomId + ":" + senderId + ":" + giftId;
Long currentCount = redis.increment(countKey, msg.giftCount);
redis.expire(countKey, 3, TimeUnit.SECONDS);
// 判断是否达到合批阈值或时间窗口结束
if (currentCount >= 10 || isWindowExpired(countKey)) {
// 推送合批消息
pushBatchedGift(roomId, senderId, giftId, currentCount);
redis.delete(countKey);
}
}
} finally {
lock.unlock();
}关键点:
锁粒度:roomId + senderId + giftId,只锁同一用户对同一礼物的连击;
窗口计数:用 Redis INCR 原子累加,EXPIRE 设置 3 秒过期;
阈值触发:累计到 10 个或窗口过期,立即推送合批消息。
Redisson vs 手写 Lua 脚本
Redisson 的 RLock 底层用 Lua 保证原子性,自带看门狗机制(锁续期)。如果业务简单,也可以手写 Lua 脚本:
if redis.call("exists", KEYS[1]) == 0 then
redis.call("set", KEYS[1], ARGV[1], "PX", ARGV[2])
return 1
else
return 0
end但要自己处理锁续期和释放时的 owner 校验。
五、礼物消息的顺序性保证
直播间礼物消息必须按发送顺序到达客户端,否则后送的礼物先显示,体验会很怪。
Kafka 分区顺序
Kafka 保证同一分区内的消息顺序。把同一直播间的礼物消息路由到同一分区:
ProducerRecord<String, GiftMsg> record = new ProducerRecord<>(
"gift_topic",
msg.roomId, // key = roomId,同一 room 进同一分区
msg
);
producer.send(record);客户端乱序检测与重排
网络传输可能乱序,客户端维护 expectedSeq,收到消息时检查:
if (msg.seq == expectedSeq) {
play(msg);
expectedSeq++;
// 检查缓冲区有没有后续消息
while (buffer.containsKey(expectedSeq)) {
play(buffer.remove(expectedSeq));
expectedSeq++;
}
} else if (msg.seq > expectedSeq) {
buffer.put(msg.seq, msg); // 暂存
} else {
// msg.seq < expectedSeq,重复消息,丢弃
}缓冲区设置大小上限(如 50 条),超过后强制播放最旧的,避免内存溢出。
六、礼物消息的过期与丢弃
用户切换直播间或断线重连后,不应该收到几分钟前的历史礼物消息。
服务端过期检查
推送前检查消息时间戳:
long now = System.currentTimeMillis();
if (now - msg.timestamp > 10_000) { // 10 秒
return; // 消息过期,丢弃
}客户端重连后的 seq 对齐
客户端断线重连后,上报最后收到的 lastSeq,服务端只推送 seq > lastSeq 的消息。如果服务端没有保存历史消息(只推实时流),直接返回当前 seq,客户端从此开始接收。
七、分布式事务:扣费与推送的一致性

用户送礼涉及两步:扣费(写订单表)和推送消息。如果扣费成功、推送失败,用户钱扣了但主播没看到。
方案一:本地消息表 + 定时补偿
扣费时同时插入本地消息表:
BEGIN;
INSERT INTO gift_order (...);
INSERT INTO outbox_message (msg_id, payload, status) VALUES (?, ?, 'PENDING');
COMMIT;推送服务轮询 outbox_message 表,状态为 PENDING 的消息重新推送,成功后标记 SENT。
方案二:事务消息(RocketMQ / Pulsar)
RocketMQ 支持事务消息,先发半消息(不可见),本地事务提交后再 commit 半消息:
TransactionSendResult result = producer.sendMessageInTransaction(msg, localTx);如果本地事务回滚,半消息自动删除,保证一致性。
总结
直播礼物后端架构的六个核心问题:
消息推送:长连接网关 + MQ + 推送服务,低延迟高并发;
去重:客户端 requestId + Redis SET NX + 数据库唯一索引;
幂等:客户端 playedMsgIds 集合 + 服务端推送去重表;
连击合批:分布式锁 + 时间窗口计数,避免刷屏;
顺序性:Kafka 分区路由 + 客户端 seq 检查与重排;
一致性:本地消息表 + 定时补偿 或 事务消息。
最后更新:
