单用户通知发送Rate Limit需求:Queue消息消费实现方案咨询
队列消息限流发送方案
核心设计思路
针对单用户速率限制、消息顺序保障和非阻塞要求,核心是用户级消息隔离+异步延迟调度:
- 为每个用户维护独立的FIFO消息队列,确保消息顺序与原流一致;
- 用调度器跟踪每个用户的发送间隔,异步触发延迟发送,不阻塞消费主线程。
具体实现步骤
- 消息分流:consumer从主队列读取消息后,按用户ID将消息投递到对应用户的专属缓存队列(比如用
HashMap<用户ID, Queue>存储,或借助Redis的List结构); - 速率检查与调度:
- 当用户队列有新消息加入时,先检查该用户上一次发送的时间戳;
- 如果当前时间与上次发送时间差≥X秒,直接发送队首消息,更新上次发送时间;
- 如果时间差不足X秒,计算剩余延迟时间,提交一个异步延迟任务(比如用
ScheduledExecutorService或异步框架的延迟API),到点后执行发送;
- 链式调度:延迟任务执行发送后,检查用户队列是否还有剩余消息,若有则再次计算下一次发送时间,提交新的延迟任务,直到队列清空。
关键细节
- 避免重复调度:每个用户同一时间只保留一个待执行的延迟任务,防止多任务触发导致速率超限(可以用一个HashMap记录用户是否有待执行任务);
- 异步选型:
- 后端代码层面:Java用
ScheduledExecutorService,Python用asyncio.create_task配合asyncio.sleep; - 中间件层面:可以用RabbitMQ延迟插件、Redis ZSet实现的延迟队列,将用户消息的下一次发送时间作为score,定时扫描触发;
- 后端代码层面:Java用
- 顺序保障:用户专属队列必须严格遵循FIFO,禁止插队发送,确保用户收到的消息顺序和原消息流一致。
示例流程(每10秒1条限制)
消息流含100条消息,用户Y有5条顺序消息M1-M5:
- 读取M1,用户Y无发送记录,立即发送M1,记录发送时间T0;
- 读取M2,此时距T0不足10秒,提交延迟任务,10秒后触发;
- M3-M5依次加入用户Y的缓存队列;
- 延迟任务触发,发送M2,更新发送时间为T0+10秒,检查到队列还有M3-M5,提交下一个10秒后的延迟任务;
- 重复上述步骤,直到5条消息按每10秒1条的频率发送完成。
内容的提问来源于stack exchange,提问作者Atul Dewangan
相关产品推荐
相关产品推荐

