直播礼物的后端架构:消息推送、去重、幂等与分布式锁

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

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


一、礼物消息推送架构

e684d1fe66e04d66aeb67d6ccbe39a6f.png

礼物特效示例:古堡秘闻

用户 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 秒过期,覆盖直播间消息的典型生命周期。


四、分布式锁:礼物连击合批

d8e78c40ded846df949db3552d220b12.png

礼物特效示例:云鲸入梦

用户在 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,客户端从此开始接收。


七、分布式事务:扣费与推送的一致性

8a8598f2d4634122abac75c6343780c0.png

礼物特效示例:星光公主

用户送礼涉及两步:扣费(写订单表)和推送消息。如果扣费成功、推送失败,用户钱扣了但主播没看到。

方案一:本地消息表 + 定时补偿

扣费时同时插入本地消息表:

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 检查与重排;

一致性:本地消息表 + 定时补偿 或 事务消息。

最后更新: