在企业级微信营销、社群运营或公告通知场景中,往往需要向成百上千个好友或微信群批量发送文本、图文链接。如果采用同步串行循环发送的方式,不仅耗时极长,还容易触发风控导致IP被封或请求超时。因此,构建一个支持限流、重试和断点续传的群发任务调度系统至关重要。

技术方案设计

本方案采用生产者-消费者模型,配合令牌桶算法进行流量控制:

  1. 任务切分:将大批量的群发名单拆分为细颗粒度的子任务。

  2. 流控限制:控制每秒请求数(QPS),避免触发平台的频率限制。

  3. 状态追踪:实时记录每个任务的发送成功、失败及待重试状态。

代码实现

以下是一个基于Node.js的异步群发任务调度实现示例,展示了如何通过控制并发节奏来安全地调用发送接口。

const axios = require('axios');

// 模拟集成平台的API配置
const API_CONFIG = {
    baseURL: 'https://www.wkteam.cn/api/v1',
    accessToken: 'your_access_token_here',
    timeout: 5000
};

/**
 * 延迟函数,用于控制请求频率
 * @param {number} ms 毫秒数
 */
const sleep = (ms) => new Promise(resolve => setTimeout(resolve, ms));

/**
 * 单个发送执行函数
 * @param {string} wxid 接收方ID
 * @param {string} content 发送内容
 */
async function sendSingleMessage(wxid, content) {
    try {
        const response = await axios.post(`${API_CONFIG.baseURL}/message/send`, {
            wxid: wxid,
            content: content
        }, {
            headers: {
                'Authorization': `Bearer ${API_CONFIG.accessToken}`,
                'Content-Type': 'application/json'
            },
            timeout: API_CONFIG.timeout
        });

        if (response.data && response.data.code === 200) {
            return { success: true, wxid };
        } else {
            return { success: false, wxid, reason: response.data.message };
        }
    } catch (error) {
        return { success: false, wxid, reason: error.message };
    }
}

/**
 * 批量群发任务调度器
 * @param {Array<string>} receiverList 接收者列表
 * @param {string} messageContent 发送内容
 * @param {number} qps 每秒发送控制数
 */
async function batchSendScheduler(receiverList, messageContent, qps = 5) {
    const intervalMs = Math.floor(1000 / qps);
    const results = { success: [], failed: [] };

    console.log(`[Scheduler] 开始执行群发任务,总数: ${receiverList.length}, 预设QPS: ${qps}`);

    for (let i = 0; i < receiverList.length; i++) {
        const wxid = receiverList[i];
        
        // 执行发送
        const result = await sendSingleMessage(wxid, messageContent);
        if (result.success) {
            results.success.push(wxid);
            console.log(`[Success] 消息已发送至: ${wxid}`);
        } else {
            results.failed.push({ wxid, reason: result.reason });
            console.error(`[Failed] 发送至 ${wxid} 失败: ${result.reason}`);
        }

        // 控频等待,防止请求过快
        if (i < receiverList.length - 1) {
            await sleep(intervalMs);
        }
    }

    console.log(`[Scheduler] 任务完成。成功: ${results.success.length}, 失败: ${results.failed.length}`);
    return results;
}

// 模拟调用测试
// const recipients = ["wxid_xxxx1", "wxid_xxxx2", "wxid_xxxx3"];
// batchSendScheduler(recipients, "这是一条系统自动群发测试消息", 2);

运营避坑指南

  1. 动态退避:当接口返回频率限制(如错误码429或特定风控提示)时,调度器应自动触发指数退避算法,暂停发送并告警。

  2. 内容变量化:群发内容建议加入随机前缀或昵称变量,降低被社交平台判定为垃圾营销机器人的概率。

Logo

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

更多推荐