微信机器人制作:消息网关的代码骨架

做微信机器人,第一个要搭的不是业务功能,而是消息网关。网关是所有微信操作的统一出入口,负责鉴权、限流、重试、路由、日志这五件事。业务代码只管调网关,不直接碰WTAPI接口。这样接口变更、Token轮换、限流策略调整,都只改网关一处。

参考API文档
WTAPI框架开发文档weiti.apifox.cn 为准。

网关的五层职责

  1. 统一鉴权:封装双Token,业务方不传Token
  2. 限流控制:发送侧令牌桶,防止触发风控
  3. 重试兜底:网络抖动自动重试,非业务错误才重试
  4. 实例路由:按策略选instanceId
  5. 日志埋点:记录每次调用的入参、出参、耗时

发送侧代码骨架

import time
import requests
from collections import deque

BASE = "https://wx.chuapi.com"

class TokenBucket:
    """令牌桶限流"""
    def __init__(self, rate, capacity):
        self.rate = rate          # 每秒生成令牌数
        self.capacity = capacity  # 桶容量
        self.tokens = capacity
        self.last = time.time()

    def acquire(self):
        now = time.time()
        self.tokens = min(
            self.capacity,
            self.tokens + (now - self.last) * self.rate
        )
        self.last = now
        if self.tokens >= 1:
            self.tokens -= 1
            return True
        return False

class MessageGateway:
    def __init__(self, token, bearer, app_id, instances):
        self.token = token
        self.bearer = bearer
        self.app_id = app_id
        self.instances = instances  # instanceId列表
        self.limiter = TokenBucket(rate=2, capacity=10)
        self._idx = 0

    def _headers(self):
        return {
            "X-finder-TOKEN": self.token,
            "Authorization": f"Bearer {self.bearer}",
            "Content-Type": "application/json",
        }

    def _pick_instance(self):
        """轮询选实例"""
        inst = self.instances[self._idx % len(self.instances)]
        self._idx += 1
        return inst

    def send_text(self, to_wxid, content, instance_id=None):
        if not self.limiter.acquire():
            raise Exception("限流中,请稍后重试")
        if instance_id is None:
            instance_id = self._pick_instance()

        body = {
            "appId": self.app_id,
            "instanceId": instance_id,
            "toWxid": to_wxid,
            "content": content,
        }
        # 最多重试3次,仅网络错误重试
        for i in range(3):
            try:
                resp = requests.post(
                    f"{BASE}/finder/v2/api/postText",
                    headers=self._headers(),
                    json=body,
                    timeout=10,
                ).json()
                if resp.get("code") == "1000":
                    return resp
                # 业务错误不重试
                return resp
            except requests.RequestException:
                if i == 2:
                    raise
                time.sleep(1)

# 使用
gateway = MessageGateway(
    token="<token>", bearer="<bearer>",
    app_id="<appId>", instances=["<inst1>", "<inst2>"],
)
gateway.send_text("<to_wxid>", "你好")

接收侧的解耦

接收侧不要在Webhook回调里直接做业务处理,要把消息丢进消息队列(如Redis、Kafka),由worker异步消费。回调接口只做三件事:签名校验、消息落库(本地消息表)、投递队列,然后立刻返回。这样即使业务处理慢,也不会阻塞回调,避免WTAPI重投风暴。

网关的工程价值

把所有WTAPI调用收敛到网关,后续要加限流策略、换Token、加监控、做多实例路由,都只改网关。业务代码永远是 gateway.send_text(...),不用关心底层细节。这是从Demo走向生产的第一个架构动作。


Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐