企业微信批量操作写成一个 for 循环,第一次跑就会撞上三件事:中间某条失败了、跑到一半进程被重启了、以及有人问「现在到哪了」。这三件都不是加个 try-catch 能解决的,它们要求你把批量任务当成一个有状态的作业来建模,而不是一段脚本。这篇讲三层模型怎么切、调用速率怎么在批量和在线之间分、断点续跑该恢复到哪一级,落地形态按 wecomapi 的调用方式来讲。
三层模型:作业、分片、条目
能续跑的批量任务至少要有三层实体,少一层就会在某个场景下失灵。
- 作业:一次提交的整体。记录参数快照、总量、状态、发起人和起止时刻。它是运营方看得懂的那一层。
- 分片:把作业切成可以并行、可以独立重跑的段。每个分片有自己的进度游标、速率闸门和状态。
- 条目:单个待执行对象。每条独立状态 —— 待执行、已提交、已确认、已失败、已放弃 —— 并记录最后一次的失败原因。
少了分片层,作业只能整体重跑或者逐条重跑,中间没有可用的粒度;少了条目层,你永远说不清「失败的三十七条到底是哪三十七条」。这两层都不能省,哪怕作业只有一百条 —— 建模成本是一张主表加一张明细表,在第一次「跑到八成挂了」的时候就回本。
条目状态里最容易被合并掉的是「已提交」和「已确认」。接口返回成功只说明请求被接受,动作是否真的生效常常要靠后续事件或事后对账才能确认。把这两个状态并成一个「成功」,续跑时就区分不了「不用重发」和「需要核实」,而这两者的处置方式正好相反。
三层之外还要一个作业级幂等键。运营在页面上双击一次提交按钮,或者上游系统重试了一次创建请求,没有这个键就是两个作业同时跑同一批人。
分片切在哪:跟着约束的作用域走
分片不是为了并行才做的,是为了让速率、失败和重跑都有一个自然的边界。所以分片键的第一原则是和约束的作用域对齐。用 wecomapi 这类统一接口执行时,频率约束和登录态通常都挂在账号维度上,那分片就先按账号切。
按账号切之后,同一账号内的条目串行或低并发执行,账号之间并行。这样速率闸门只需要在分片内部生效,不用跨分片协调;进程重启后各分片各自恢复,也不会出现两个分片互相抢同一个账号速率的情况。账号内还要不要继续切,看单账号条目数会不会大到一个分片跑几小时 —— 会的话按对象标识哈希再切一层,注意并发度仍受账号约束,这次切分只是为了缩短单次重跑的范围。复合分片键的具体推导站内的群发那篇讲得更细,这里只讲分片作为调度实体的部分。
分片一旦创建,它就是一个有生命周期的对象,而不是一个下标区间。它至少要维护三样东西:进度游标(下一条从哪开始)、状态(待跑、跑着、暂停、已完成、已失败)、最后推进时刻。第三样最容易被漏掉,而它是唯一能发现「分片卡死」的信号 —— 分片卡死的表现不是失败率上升,是这个数字不动了,错误率一片干净。
- 分片数在作业创建时固定,运行中不要调整 —— 分片数一变,已完成的游标就对不上了,续跑时会有条目被跳过或者跑两遍;要提高并发就调 worker 数量,别动分片本身。
- 分片之间不共享可变状态。要看总进度就各自上报,由作业层汇总。
有一种分片方式几乎总是错的:按提交顺序等分,前一千条一片、后一千条一片。它会把同一个账号的条目撒进所有分片,速率闸门立刻需要跨分片协调,而且失败时你没法只重跑某个账号。
限速:闸门开在分片里,预算记在作业上
批量任务的限速和在线请求的限速不是一回事。在线请求求的是「别被拒」,批量任务还得求「别把在线请求挤掉」—— 两者共用同一个账号的发送节奏,而客服消息晚十秒是投诉,批量触达晚十分钟没人在意。
所以批量任务的速率不能设成「能跑多快跑多快」。工程上是两个独立的限速器:在线链路一个,批量任务一个,批量那个的上限等于你给这个账号定的总速率减去在线保留的水位,而不是动态去抢。抢的结果永远是批量赢,因为它一直在跑。
// 示意逻辑,精确字段与端点以线上文档为准
async function runShard(shard) {
for await (const item of shard.pending()) { // 只取未完成的条目
if (shard.cancelled) break; // 取消在条目边界生效
await batchLimiter.acquire(shard.accountId); // 批量专用限速器,不与在线抢
await store.mark(item.id, "submitted"); // 先落状态,再发请求
try {
await post("https://manager.wecomapi.com/message/sendText", {
guid: shard.guid,
toId: item.toId,
content: render(item),
});
await store.mark(item.id, "confirmed");
} catch (e) {
await store.fail(item.id, classify(e)); // 失败分类要留,重跑靠它筛
}
await shard.advance(item.id); // 推进游标,决定能不能续跑
}
}- 条目之间留随机间隔,别用固定 sleep。固定间隔只是把一个洪峰切成几个等距的小洪峰。
- 速率要能在作业运行中热调,而不是只能停掉重提。放量观察本来就是分阶段进行的。
- 给作业设「总量预算」而不只是「速率」:这一次最多发出多少条,超了自动暂停等人确认。速率管的是瞬时,预算管的是总量,两个都需要。
- 被拒之后别在分片内原地死等。把这条退回待执行、让整个分片降速,让位给还能跑的条目。
退避重试的写法和错误分类站内另有一篇专讲,这里只强调批量场景下的优先级:降速排在重试前面。重试是在同样的节奏上再试一次,如果节奏本身就是问题,重试只会让它更糟。
断点续跑:恢复到哪一级
「续跑」有三种粒度,多数事故是因为只实现了最粗的那一种。
- 1作业级:整个作业重来,靠条目状态过滤掉已完成的。实现最省事,问题在于如果作业参数引用的是「当前符合条件的客户」,重跑时人群已经变了。
- 2分片级:只重跑失败或未跑完的分片。适合分片彼此独立的场景,这是默认应该做到的粒度。
- 3条目级:只重跑失败的条目,其余一律跳过。这才是真正意义上的续跑,前提是每条都落了状态,而且落状态发生在发请求之前。
决定能不能做到条目级的,就是写状态和发请求的顺序。先发请求再记状态,进程在这两步之间被杀掉,这一条就永远停在「不知道发没发」;先记「已提交」再发,最坏情况是有一条被跳过,而它能在事后对账里被捞出来。两害相权,宁可漏发一条也不要重发一条 —— 前者能补,后者补不回来。
另一件必须在作业创建时就定死的是目标集合怎么固化。按查询条件跑的作业,续跑时到底用原来那批对象,还是重新查一次?两种都有合理场景,但必须显式选:补发通知类的用快照(就是原来那批),状态同步类的用重查(当前那批)。不选的后果是它由「作业重启那一刻数据库里是什么」随机决定,而这种不确定性在复盘时几乎无法还原。
落地上还有一条容易被忽略:条目状态表是你自己的。wecomapi 侧只负责接受一次调用,「这一条到底走到哪一步」必须由你这边记录,不要指望回查能替你把中间状态还原出来。
- 用快照就在作业创建时把对象列表落成条目,之后只认这张表。
- 用重查就把查询条件和游标一起存,续跑时从游标继续,并接受集合可能已经变化。
- 两种都要记录作业创建时刻。排查「为什么这个人不在名单里」,靠的就是这个时刻。
暂停、取消,以及跑歪了怎么办
能停下来是批量作业和脚本最大的区别。三个动作要分清语义并写进说明,因为运营在紧急时刻按下去的那一下,需要确定它到底做了什么。
- 暂停:停止取新条目,已提交未确认的继续等结果,可恢复。这是最常用的一个,也是唯一应该做成一键的。
- 取消:停止取新条目并把剩余条目置为已放弃,不可恢复。取消不撤销已经做过的动作,界面上必须写清楚这一点。
- 终止:进程级停止,可能留下一批「已提交未确认」的条目,恢复之前必须先对账。
取消要在条目边界生效,不能在一条动作执行到一半时打断。实现上就是循环体开头检查一次标志位,而不是去杀线程 —— 杀线程留下的中间态,正是最难对账的那一类。
作业层还要有熔断:连续失败超过阈值,或失败率超过某个比例时自动暂停并告警。批量任务最坏的失败模式不是跑失败,是带着一个写错的模板把两万条全都成功发出去了。熔断挡不住内容错误,但能挡住「凭证过期导致全部失败还在傻跑一整夜」这一类,而后者的发生频率高得多。
- 1随机杀掉一个执行进程,作业能不能自己恢复,恢复之后有没有重复条目。
- 2把某个账号的调用人为全部打成失败,作业会不会拖着不动,熔断有没有按预期触发。
- 3同一个作业提交两次,会不会真的执行两遍。
- 4运营在界面上能不能看懂「跑到哪了、失败多少、失败原因分几类」。看不懂就会有人在旁边手工再做一遍。
- 5作业跑完之后,「已提交未确认」的那些条目有没有人管、多久之内会被核实。
本文讲的是调度模型与工程取舍。精确的字段、端点与调用约束以 wecomapi 文档为准,示意代码不要照抄上生产。
常见问题
- 批量任务和实时消息能共用一套服务吗?
- 代码可以共用,节奏和通道不能。用 wecomapi 时两条链路走的是同一个账号实例,账号侧的发送节奏是共享的:批量迟到十分钟没人发现,客服消息晚十秒就是投诉,两者一起抢这个节奏时永远是批量赢,因为它一直在跑。做法是给批量任务单独的限速器和更低的优先级,上限取你为该账号定的总速率减去在线保留水位之后的那部分,两条链路的积压指标也要分开看。
- 作业跑到一半改了模板或名单,已经跑的那部分算哪一版?
- 算作业创建时那一版,前提是你在创建时对参数和对象列表做了快照。这也是参数不该在运行中热改的原因 —— 要改就取消当前作业、按新参数重开一个,让「这一批用的是哪个版本」永远只有一个答案。中途改参数省下的那点时间,远不够抵消事后对不上账的排查成本。
- 失败的条目该自动重跑还是等人?
- 按失败分类分开处理。限速类和网络超时类退避之后自动重跑;参数类、权限类自动重跑纯属浪费,直接标记已放弃交给人;分类不出来的先进人工队列,别默认重试。自动重跑还要有次数上限和总时长预算,否则一次凭证问题能让作业空转一整天。具体的错误分类以 wecomapi 文档为准,不要按二手描述硬编码。
准备好动手了?
精确字段、鉴权与端点以线上文档为准;可在控制台创建密钥后联调。
