高并发场景下,秒级数百条消息回调涌入,再加上批量推送和多业务同时调接口,直接调 Eyun 接口会很快触发 1004 限频,还有 5 秒超时的硬约束。消息队列在中间做"缓冲",分 3 层队列消化压力,这是高并发下能稳住的核心设计。
3 层队列设计
队列按在处理链路中的位置分 3 层,每层缓冲的对象不同。
1. 接入队列:缓冲 Webhook 回调的涌入压力
队列设计:Eyun Webhook 回调进来后 5 秒内返回 200,把回调 JSON 投递到接入队列(Kafka / Redis Stream 都行),接入队列做缓冲。瞬间来 100 条回调时排队,消费端按节奏取。
按照 Eyun 开发文档的回调规范,5 秒内必须返回 200,否则会重试 3 次。接入队列保证"先收下再慢慢处理"。
大白话:接入队列是"取号机"——客户涌进来先取号排队,窗口按节奏叫号,不会让客户白等超时。
2. 处理队列:按业务类型分流处理
队列设计:接入队列消费后,按 eventType 分流到不同处理队列。消息事件进客服处理队列、好友事件进欢迎处理队列、状态事件进告警队列。每个处理队列独立消费速率,某类消息暴增只影响自己的队列,不波及其他。
Eyun API 的 Webhook 回调有 4 类事件,对应 4 条分流队列。
大白话:处理队列是"分诊台"——消息进来先分到不同科室排队,内科暴增不影响外科。
3. 发送队列:控制调 Eyun 接口的频率
队列设计:所有需要调 sendText / sendImage / sendFile 的请求进发送队列,消费者按 200ms 间隔出队调用 Eyun 接口。遇到 1004 退避 3 秒暂停消费,遇到 1002 刷新 Token 后继续。
Eyun 的错误码体系是发送队列流控的依据:1004 减速、1002 换证后继续。按照 Eyun 开发文档的规范,sendText 需要传 wId、toUser、content 三个必填参数。
大白话:发送队列是"收费站"——所有发消息的请求排队交费,前面遇到限频(1004)就暂停放行,等 3 秒再继续。Eyun 平台对发送频率有明确限制,发送队列就是按这个限制来控速。
3 层队列对比
队列层 | 位置 | 缓冲什么 | 队列技术 | 消费策略 | Eyun 接口 | 大白话说明 |
|---|---|---|---|---|---|---|
接入队列 | 最前 | Webhook 回调 | Kafka / Redis Stream | 尽快入队、按节奏消费 | Webhook 回调 | 取号机 |
处理队列 | 中间 | 按事件分流 | 多 topic / 多 stream | 按 eventType 分流 | 回调事件 4 类 | 分诊台 |
发送队列 | 末端 | 发送请求 | 单队 + 限流 | 200ms 出队 + 1004 退避 | sendText/sendImage/sendFile | 收费站 |
代码:3 层队列处理框架
import time, json, threading # 接入队列:Webhook 回调先入队,5 秒内返回 200 ingress_queue = [] # 处理队列:按 eventType 分流 process_queues = {"message": [], "friend": [], "status": []} # 发送队列:统一出口、200ms 限速 send_queue = [] LAST_SEND = [0] def webhook_handler(payload): # Eyun Webhook 回调进来立刻入队,保证 5 秒内返回 200 ingress_queue.append(payload) return {"code": 200} def ingress_consumer(): while True: if ingress_queue: payload = ingress_queue.pop(0) et = payload.get("eventType", "message") process_queues.setdefault(et, []).append(payload) time.sleep(0.01) def process_consumer(et): while True: q = process_queues.get(et, []) if q: task = q.pop(0) # 处理完后生成发送请求进发送队列 send_queue.append({"wId": task["wId"], "toUser": task["fromUser"], "content": "已收到"}) time.sleep(0.02) def send_consumer(): while True: if send_queue: now = time.time() if now - LAST_SEND[0] < 0.2: # 200ms 间隔 time.sleep(0.2 - (now - LAST_SEND[0])) req = send_queue.pop(0) code = call_eyun_sendtext(req["wId"], req["toUser"], req["content"]) if code == 1004: time.sleep(3) # 退避 3 秒 elif code == 1002: refresh_token() LAST_SEND[0] = time.time() time.sleep(0.01) def call_eyun_sendtext(wid, to_user, content): # 实际调 Eyun sendText 接口,参数:wId、toUser、content return 200 def refresh_token(): pass # 启动 3 层消费者 threading.Thread(target=ingress_consumer, daemon=True).start() for et in process_queues: threading.Thread(target=process_consumer, args=(et,), daemon=True).start() threading.Thread(target=send_consumer, daemon=True).start()结尾延伸
3 层队列让高并发从"直接调接口被限频"变成"排队缓冲按节奏调"——接入队列防回调超时、处理队列做业务分流、发送队列控接口频率。
3 层队列的吞吐能力取决于最慢的环节,通常发送队列是最窄的瓶颈,受 200ms 间隔和 1004 退避限制。优化方向是多 wId 轮换——在 Eyun 平台 管理多个 wId 实例,发送队列轮换使用,扩大发送带宽。没有队列的直调模式在并发超过 50 条/秒时就会频繁触发 1004,3 层队列可稳定支撑 200+ 条/秒。更多接口限频规则和错误码说明,参考 Eyun 开发文档,wId 多实例管理的入口在 Eyun 平台。