微信群发消息的Java实现:分片发送与失败重试机制详解

发布时间:2026/10/10 6:43:39
微信群发消息的Java实现:分片发送与失败重试机制详解
前阵子有个业务方找到我说要做微信群发消息一次性给好几千个客户推送通知。需求听起来很简单写个接口把消息发给所有人。可真把代码写起来才发现就这么一个“for循环调接口”的活儿坑比想象中深得多。直接循环调用发到一半被平台限流接口报错也不知道哪批失败了客户收到重复消息还跑来投诉。最后我老老实实把方案重构成了分片发送 失败重试机制才把这个问题彻底收住。这篇文章我就把完整的优化过程、代码思路和踩坑记录分享出来。如果你的项目里也涉及Java对接微信群发、短信群发或者其他类消息API接口的批量处理这套思路完全可以平移过去复用。1. 先理清核心矛盾为什么大批量发送不能直接for循环很多人拿到“群发消息”的第一反应就是写个for循环把用户列表遍历一遍循环体里面调一次发送接口。这个方案在几十条、一百条的时候勉强能跑但上了千条之后必然炸。原因倒不是Java循环慢而是你忽略了两个最要命的东西接口端的频率限制、单条发送的失败概率。1.1 接口限制与业务痛点的拆解微信群发消息这类API接口包括公众号群发、企业微信群发、模板消息推送等几乎都有明确的频率限制。以常见的接口规范为例有的按分钟维度限制调用次数有的按单次请求的最大接收人数限制批量规模还有的会针对同一内容做去重校验。你的循环体调得再快到了接口这一层一样会被拦住返回限流错误码甚至触发更严厉的临时封禁。另外还有一个很容易被忽视的事实一条消息发出去底层要经过网络传输、平台内部路由、对方设备投递。网络一抖动或者平台某个节点抽风单条消息失败是常态概率可能达到百分之一甚至更高。几千条消息就算成功率达到99%也有几十条是失败的。如果不做重试用户那边就会出现有人收到了、有人没收到的情况这在线下业务场景里几乎等于事故。我把问题总结成下面的表格这也是我当时理清需求的起点痛点直接原因后果接口限流请求频率过高报错、封禁、发送中断单条失败网络抖动、平台异常消息丢失、用户投诉重复发送重试逻辑设计不当用户反感、平台拦截状态不可知没有记录明细排查困难、无法对账1.2 整体设计思路先拆再送失败兜底解决问题的思路说穿了一点都不玄乎既然一批发不完那就分批发既然会失败那就失败了再补一次。但具体到工程实现有两个关键点决定了这套方案的上限。第一分片发送不能只按数量硬切。你得考虑接口的承载能力、发送的时间分布、以及平台对相同内容短时间内重复提交的限制。我后面会详细讲怎么确定分片大小和发送间隔这里先记住一个原则分片的目标不是“把大列表切成小列表”而是“让发送速率平稳落在接口允许的范围之内”。第二失败重试不是简单地在catch里再调一次。你必须先区分哪些失败值得重试、哪些失败重试一万次也没用还要设计重试的退避策略防止同一批失败的消息在同一个时间点集体重试直接把接口再次打爆。这两点想清楚了代码的骨架也就出来了一个任务拆成分片分片进入线程池调度发送发送结果记录状态失败的消息进入重试队列重试也走同样的分片控制超过重试上限的进入死信人工处理。整体流程图在我脑子里过一遍之后落地就是Spring Boot工程里面几个互相配合的组件。2. 分片发送机制把大任务拆成接口能承受的小块分片发送是整个优化方案的地基。地基打不好后面重试做得再漂亮也没用。我这一节会讲分片参数怎么定、并发怎么控制、以及分片发送与平台API交互时的实际代码怎么写。2.1 分片大小与发送间隔的确定分片大小的确定说实话第一版我也是拍脑袋定的后来在压测环境里调了三轮才找到合理的做法。核心参考维度有三个接口单次允许的最大接收人数、接口分钟级调用上限、以及单条消息发送的平均耗时。举个例子假设接口单次最多接收100个人那分片大小就不能超过100我习惯取80留20%的余量给参数校验、用户解析这些额外消耗。假设接口限制每分钟调用60次单次调用平均耗时300毫秒那我每秒最多发3到4个分片。如果分片大小是80每片耗时300毫秒那么理论上每秒钟能处理240个人一分钟就是14400个人这个量级对大多数微信群发场景已经足够了。真正要控制的是不同分片之间的发送间隔。代码里我直接用Thread.sleep()显然太粗暴更好的做法是用限流器来控制发送速率。我第一版用的是ScheduledThreadPoolExecutor直接把每个分片任务按固定间隔丢进去执行后来发现这个方案对“速率控制”的理解更直观代码也更简单。下面是一段核心实现public class ChunkSendScheduler { private static final int CHUNK_SIZE 80; private static final int MAX_QPS 4; // 每秒最多发4个分片留足余量 private final ScheduledExecutorService scheduler Executors.newScheduledThreadPool(2); private final MessageSendClient messageSendClient; public void scheduleSend(ListString userIds, String content) { ListListString chunks Lists.partition(userIds, CHUNK_SIZE); long delay 0L; for (ListString chunk : chunks) { scheduler.schedule(() - sendChunk(chunk, content), delay, TimeUnit.MILLISECONDS); delay 1000L / MAX_QPS; // 每个分片间隔250ms } } }这段代码做的事情很简单先把用户列表按80人一批切开然后让每个分片任务按250毫秒的间隔依次执行。用ScheduledExecutorService的好处是任务的触发节奏由调度器统一控制不会因为某个分片执行慢了导致后面的请求瞬间堵在一起。2.2 发送节奏控制的进阶令牌桶思路如果你对接的API接口限制更复杂比如同时限制每秒调用数和每分钟调用数上面这个固定间隔的方案就不够用了。我在另一个项目里用过令牌桶的思路效果更稳。令牌桶说白了就是“匀速往桶里放令牌请求来了必须拿到令牌才能执行”。Guava的RateLimiter就是现成的实现。用起来代码非常简洁RateLimiter rateLimiter RateLimiter.create(4.0); // 每秒放4个令牌 public void sendWithRateLimit(ListString chunk, String content) { rateLimiter.acquire(); // 拿不到令牌就阻塞等待 messageSendClient.sendGroupMessage(chunk, content); }对比固定间隔的方案令牌桶最大的优势是能应对执行耗时的波动。比如某一次接口响应特别慢耗时从300毫秒涨到了800毫秒固定间隔方案里后续任务依然按原节奏触发会造成任务在线程池里堆积令牌桶方案则会让acquire阻塞更久实现动态拉长间隔避免请求积压。实测下来接口响应波动明显的场景令牌桶成功率比固定间隔高出不少。2.3 分片任务线程池的配置细节分片任务不能直接丢进单线程的调度器里发不然一个分片失败重试会堵住后面的所有发送。我单独建了一个发送线程池调度器只负责任务的触发真正的发送逻辑在线程池里执行。线程池参数我踩过一次坑。最开始照着网上推荐的CPU密集/IO密集公式配corePoolSize设8maxPoolSize设20队列容量塞了2000。结果高峰期所有分片任务全堆在队列里发送延迟越来越大最后接口限流错误和超时错误一起涌出来。后来我意识到了一个问题这种场景下你根本不需要那么大的线程池因为接口的QPS上限就摆在那里线程再多也只会加重无效请求。我这里最终确定的是corePoolSize和maxPoolSize都设成4到6队列容量设成分片总数的一半。为什么是这个数因为前面限流已经限制了每秒最多4个分片4到6个线程已经足够消化这些任务线程再多反而会引入不必要的上下文切换。队列容量控制在一半则是有意为之任务太多直接走拒绝策略往数据库落状态而不是无限堆在内存里等。这是“宁可拒绝不可积压”的思路后面篇幅我会再展开说。2.4 发送接口调用的落地实现分片发送最后落到API接口调用时有几个细节很值得注意。我先贴一段实际发送的代码再逐一解释public ChunkSendResult sendChunk(ListString userIds, String content) { long start System.currentTimeMillis(); int retryCount 0; while (retryCount MAX_RETRY) { try { SendResponse resp messageSendClient.sendGroupMessage(userIds, content); if (resp.isSuccess()) { return ChunkSendResult.success(userIds.size(), System.currentTimeMillis() - start); } if (!resp.isRetryable()) { return ChunkSendResult.fail(userIds, resp.getErrorMsg(), false); } retryCount; long waitMs computeBackoffTime(retryCount); Thread.sleep(waitMs); } catch (InterruptedException e) { Thread.currentThread().interrupt(); return ChunkSendResult.fail(userIds, interrupted, true); } catch (Exception e) { retryCount; long waitMs computeBackoffTime(retryCount); try { Thread.sleep(waitMs); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); return ChunkSendResult.fail(userIds, interrupted, true); } } } return ChunkSendResult.fail(userIds, exhausted, true); }发送异常这块最重要的是区分可重试与不可重试。我遇到的不可重试情况主要有三种参数格式错误比如用户ID传成了null、权限错误比如appSecret过期、内容违规被拦截。这三种错误重试多少次结果都一样白白消耗资源还有可能把账号搞进更严厉的风控名单。可重试的则包括网络超时、服务端5xx错误、限流错误返回特定的code、以及一些临时性的系统异常。判断逻辑我用响应对象里的一个字段来标识这比在catch里猜异常来源更精确。服务端明确告诉你“限流了”跟你自己网络超时处理策略应该是不一样的限流重试要等待更长的时间网络超时则可以相对快速地重试。2.5 分片状态跟踪每个分片发送完我当时都会往数据库更新一条记录。字段不多但非常关键分片ID、任务ID、分片序号、接收人数、成功人数、失败人数、状态、耗时、错误信息、重试次数。为什么要记录分片级别的状态而不是只记录总状态原因很实在如果整个任务半路崩了重启之后我需要知道哪个分片发出去了、哪个没发出去、哪个发了一半才能决定是继续还是重新发。没有这些明细唯一的选择就是把整个任务再跑一遍重复发送的风险直接拉满。数据库表结构我当时是这么设计的可以参考CREATE TABLE send_chunk_record ( id BIGINT PRIMARY KEY AUTO_INCREMENT, task_id VARCHAR(32) NOT NULL, chunk_index INT NOT NULL, total_count INT NOT NULL, success_count INT DEFAULT 0, fail_count INT DEFAULT 0, status TINYINT NOT NULL COMMENT 0待发送 1发送中 2成功 3失败待重试 4最终失败, retry_times INT DEFAULT 0, error_msg VARCHAR(500), cost_ms BIGINT DEFAULT 0, create_time DATETIME NOT NULL, update_time DATETIME NOT NULL, UNIQUE KEY uk_task_chunk (task_id, chunk_index) ) COMMENT 分片发送记录表;这里的唯一索引uk_task_chunk是防止重复的关键。后续不管是任务重启还是定时补偿只要往这张表插数据重复的分片自然会被数据库挡下来这比在代码里用各种状态判断要可靠得多。3. 失败重试机制让发送失败的消息有机会被补救如果说分片发送解决了“发得太快”的问题那失败重试解决的就是“发了没收到”的问题。这一节我重点讲重试策略的选择、重试队列的设计、以及幂等性怎么保证。3.1 可重试与不可重试的区分逻辑前面发送代码里区分了可重试和不可重试现在展开讲讲我的判断依据。总的思路就是一句话重试要针对那些“这次失败但下次可能成功”的情况。我遇到过的情况可以分为三类。第一类是平台返回明确错误码的比如“当前调用过于频繁”“服务内部错误”“系统繁忙”这些都是临时性的过一会儿再试大概率能成功属于可重试。第二类是网络层面的异常比如SocketTimeoutException、ConnectException网络抖动恢复之后请求就能成功也属于可重试。第三类是业务层面的错误比如某个用户ID无效、内容被判定违规、没有群发权限这些属于永久失败重试没有任何意义甚至可能因为反复提交同样的内容触发更严重的账号风控。判断逻辑放在哪个环节也很重要。我推荐在发送客户端就完成判断并且把“是否可重试”作为响应对象的一个属性返回而不是在重试框架里靠catch异常类型去猜。原因很简单HTTP状态码200不代表发送成功400也不代表绝对不能重试。只有真正解析过API响应结构的人才知道哪些code背后到底是什么含义。3.2 重试退避策略指数退避加随机抖动重试策略里最忌讳的就是“失败后立刻重试、再失败立刻再重试”。几十分片同时失败如果全部立刻重试接口瞬间收到同样数量的请求等于再触发一次限流形成恶性循环。正确的做法是加上退避时间让每次重试的间隔越来越长。我用的公式是private long computeBackoffTime(int retryCount) { long base Math.min(30000L, 1000L * (1L Math.min(retryCount, 5))); long jitter ThreadLocalRandom.current().nextLong(0L, Math.max(1L, base / 5)); return base jitter; }这里的逻辑是第一次重试前等1秒左右第二次等2秒左右第三次4秒以此类推最多封顶30秒。jitter是随机抖动目的是让同一批失败的分片不会在同一个毫秒级时间点上集体重试。这个抖动在别的场景里可能无所谓但在批量重试场景里非常关键——系统设计上有一个“惊群效应”几十上百个任务同时醒来打接口跟定时炸弹没什么区别。实际运行效果是这样第一轮重试大约1.2秒后开始第二轮大约2.5秒第三轮大约4.8秒。整体重试节奏平稳接口压力被自然地摊开了。3.3 最大重试次数与消息生命周期重试不能无限做下去。我最终定的策略是单条消息最多重试5次超过5次就进入死信状态。这个数字不是拍脑袋定的而是根据业务容忍度和接口恢复时间综合算出来的如果接口连续5次都失败大概率不是临时抖动而是持续性的故障你再重试也只是增加无效请求不如先把数据保住等人工介入。进入死信状态的消息怎么处理我们当时有两个手段。一是定时任务扫描死信表往企业微信群里推告警由运营人员用管理后台手动重发。二是写了一个补偿接口手动选择死信记录后重新走一遍分片发送流程。这里有个细节死信消息的重新发送要重新走限流逻辑不能直接一条条硬调接口不然等于绕过了分片控制前功尽弃。3.4 重试过程中的幂等性保证做群发消息的人最怕两件事消息没发出去以及消息发了两遍。重试机制天然会带来“可能发了两遍”的问题幂等性设计必须提前做好。我的做法是给每条待发送的用户记录生成一个全局唯一的消息ID格式可以是taskId _ userId。发送请求时带上这个ID平台接口如果支持去重就会处理相同ID的重复请求如果接口不支持我在自己的数据库里用唯一索引挡住重复提交。数据库层面的设计是这样的发送明细表里有一个消息ID字段加上唯一索引。每次重试之前先往明细表插入一条状态为“发送中”的记录如果插入时违反唯一约束说明这条消息已经处理过哪怕前一次状态没来得及更新直接跳过不再发送。这个方案在后端系统里非常常用逻辑简单效果可靠。CREATE TABLE send_message_detail ( id BIGINT PRIMARY KEY AUTO_INCREMENT, message_id VARCHAR(64) NOT NULL COMMENT 全局唯一消息ID, task_id VARCHAR(32) NOT NULL, user_id VARCHAR(64) NOT NULL, status TINYINT NOT NULL COMMENT 0发送中 1成功 2失败 3死信, retry_times INT DEFAULT 0, error_msg VARCHAR(500), create_time DATETIME NOT NULL, update_time DATETIME NOT NULL, UNIQUE KEY uk_message_id (message_id) ) COMMENT 发送明细表;每次发送动作的事务边界是先插入明细记录再调用发送接口最后更新明细状态。你把“插入动作”当成一次发送令牌的领取抢到了才允许发送抢不到就说明别人正在发或者已经发过了。这个思路比单纯依赖内存状态要稳尤其是在多个实例部署的情况下。4. 任务状态管理与监控做到对每一条消息都心里有数分片发送和失败重试这两块地基打完之后你还缺一层“统筹全局”的能力。几千条消息的批处理任务如果没有一个任务级别的状态管理体系运行到一半你根本不知道整个任务处于什么阶段更谈不上对账和排查。4.1 任务总表与状态流转设计我在分片记录表之外单独建了一张任务总表记录每个群发任务的整体信息。字段相对简单任务ID、任务名称、接收用户总数、分片总数、成功分片数、失败分片数、任务状态、开始时间、结束时间、创建人、备注。任务状态的流转我定义为待处理、处理中、部分成功、成功、失败。你可能会问“部分成功”和“失败”有什么区别区别在于失败是否还需要任务是“终态”。如果整个任务还有分片在重试中状态就是处理中所有分片都结束了但有一部分失败无法重试就是部分成功所有分片都成功才是真正意义上的成功。这套状态流转保证了业务的准确性。运营后台对外展示的时候用户看到的不是“成功”或“失败”这种一刀切而是“成功980人失败20人其中15人重试成功5人需要人工处理”这个信息对业务决策非常关键。4.2 定时补偿任务光有状态记录还不够因为系统可能在任何时候崩溃。我当时加了一个定时补偿任务每5分钟扫描一次任务总表把所有处于“处理中”但超过30分钟没有更新状态的任务捞出来再把它们对应的分片记录查出来看看卡在哪个环节。补偿的逻辑也很简单如果某个分片状态是“发送中”但超过10分钟没有变化就把它重新置为“待发送”重新丢进发送调度器。为了防止一个分片被两个线程同时处理我上面说的唯一索引这时候就发挥作用了——重复的明细记录插入会被数据库挡住不会出现真正意义上的重复发送。这个定时补偿任务可以说是整套方案的“安全带”。有了它我不再担心半夜告警响起来之后需要人工去数据库改状态系统自己会把卡住的任务拉回正轨。4.3 监控指标与日志规范代码写得再好没有监控你也只能在故障发生之后被动响应。我给这套群发系统加了三个最关键的监控指标。第一个是发送成功率成功消息数 / 总消息数。这个指标从90%掉到80%就说明接口开始不稳定了需要关注。第二个是分片平均耗时如果从300毫秒涨到1000毫秒说明接口响应变慢可能接近限流阈值。第三个是重试率重试消息数 / 总消息数。重试率超过20%大概率是分片大小或者发送间隔设置得不合理需要调参。日志方面我给自己定了一个铁律每个分片发送结束必须输出一条包含分片ID、任务ID、成功数、失败数、耗时的日志每次重试动作也必须有日志格式统一为[taskId][chunkIndex][retryTimes] errorMsg。这样排查问题的时候直接按taskId全局搜索就能完整还原这个任务的整个生命周期。4.4 Spring Boot工程里的组件划分上面讲的这些逻辑落到代码工程里我在一个服务类里做了模块化处理。核心组件如下chunkSendScheduler分片调度器负责任务拆解与发送节奏控制。messageSendClientAPI调用客户端封装了微信群发消息的HTTP调用、响应解析与可重试判断。sendRecordService分片记录与明细记录的读写服务。retryQueueHandler重试队列处理器负责从数据库中捞取待重试的分片重新推入调度器。failedMessageHandler死信处理器负责超过重试上限的消息的告警与人工补偿入口。这种划分的好处是每个组件可以独立测试。我当时在本地用Mock接口做了完整的单元测试分别验证了限流场景、网络超时场景、重复提交场景确保每个组件的行为符合预期之后再联调真实接口问题排查的范围一下就缩小了。5. 实战问题排查与经验总结任何架构设计都得经过真实流量的检验。这套方案上线后我遇到了一些很有意思的问题有的是设计阶段没想到的有的是想到了但没有做到位的。这一节我把它们整理出来相当于一份排错清单。5.1 高频问题与解决方案速查表问题现象原因解决方案跑到一半任务没有任何报错就停了线程池队列满了任务被拒绝调整队列容量增加拒绝策略的告警同一条消息收到两次定时补偿与正常发送同时处理同一批次依赖唯一索引挡住重复插入重试成功率越来越低退避间隔太短重试请求依然撞上限流加大基础退避时间增加随机抖动接口大量返回限流错误分片间隔计算错误实际QPS超过限制用RateLimiter替换固定Sleep消息状态一直是“发送中”进程崩溃状态没来得及回写增加定时补偿任务超时自动重置数据库连接池耗尽发送明细写入频繁连接不够用调大连接池明细写入改批量插入5.2 线程池与数据库连接池的联合考量这个问题花了我不小的代价才彻底想明白值得单独说说。分片发送涉及两个池线程池和数据库连接池。发送动作会调接口耗时较长这个期间线程是阻塞的但它占用的数据库连接是释放了的。而插入明细记录的动作很快但频率很高。如果这两者共用同一个连接池大批量发送时所有线程都在等待接口响应数据库连接反而被“快速插入”动作大量占用最终连接池被打满插入动作反而成为瓶颈。我的解决办法是给不同的操作单独建设连接池或者至少给明细表写入设置独立的连接池参数。发送线程池持有少量的连接用于低频状态更新明细表写入使用另一个池pool size设得更充裕一些。这个调整之后数据库连接池耗尽这个问题就再没有出现过。5.3 Spring事务与长任务不能共存还有一个问题在初版代码里出现过就是我把整个分片的发送逻辑包在了一个Transactional方法里。理由现在看来很天真担心数据不一致想用事务兜底。结果发送接口要等300毫秒网络超时要等几十秒事务迟迟不提交数据库连接一直被占用加上Spring事务默认隔离级别下的锁行为直接拖垮了数据库性能。后来我把事务去掉改用分布式状态管理。每个明细记录的插入和更新都是单独的短事务发送动作完全放在事务之外。数据不一致的风险由前面说的幂等性设计和定时补偿任务来兜底效果远好于一个大事务包住一切。5.4 测试方法与回放机制最后聊聊测试。群发消息的接口肯定是不能直接拿正式用户来测的尤其是压测的时候一条测试消息发出去用户真能收到。我当时的做法是搭了一个Mock发送服务模拟平台的限流、超时、5xx错误然后把测试数据量拉到1万人完整地跑一遍分片发送和失败重试的全流程。这套Mock服务还支持故障注入可以指定某个分片永远失败、指定某个错误码触发限流、甚至模拟进程中途宕机。每次改动上线之前我都会跑一遍这三种场景确认系统能够按照预期进行状态流转和补偿。还有一个值得说的细节真实线上出过一次问题之后我把当时的请求参数、错误码、时间点全部导出写了一个故障回放工具。用同一份数据重新跑一遍优化后的代码对比前后行为差异。这个方法帮我验证了多个修复方案的有效性建议做消息类系统的朋友也可以试试“故障回放”这个思路。最后再分享一个我的实操心得这套分片发送与失败重试机制我前后迭代了好几版从最初简单的循环加重试到现在调度器、限流器、补偿任务、幂等表完整配合稳定运行的时长已经超过一年。我个人在实际操作中最深的体会是批量处理优化本质上不是把代码写得更快而是把失败处理得更稳。Java里一个for循环跑几千次其实很快真正慢的是接口的响应等待真正不稳定的也是网络和外部接口。分片解决的是“别太快”重试解决的是“别丢消息”幂等解决的是“别重复”。把这三件事想清楚代码上的东西其实都是水到渠成的细节实现。如果你现在正好要做微信群发消息的批量接口建议第一步先把自己要对接的接口文档翻一遍把限流规则、错误码表、字段约束都搞清楚再回来对照这篇文章里的设计。先把分片和重试的主流程跑通再逐步加状态管理、监控和补偿任务。不要一上来就想着把系统做得面面俱到稳定性和可观测性是边跑边补的核心是先让批量发送这个动作变得可控、可回放、可排查。后续如果业务量继续上涨还可以考虑把重试队列从数据库搬到消息队列中间件里把分片调度的节奏配置化配合压测平台做自动化容量验证。这些扩展我在另一个项目里已经落地了一部分等有时间再专门写一篇记录。