恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
Java群发消息API批处理:分片发送与失败重试机制详解
首页
资讯中心
/
Java群发消息API批处理:分片发送与失败重试机制详解
Java群发消息API批处理:分片发送与失败重试机制详解
发布时间:2026/10/11 13:42:51
做Java后端的人大概率都碰过微信群发消息API接口的对接。你以为拿到用户列表、写个for循环逐个调接口就完事了恰恰是这种最朴素的批量处理写法在生产环境能把一批人坑得够惨。这篇文章把我在真实项目里用到的Java批量处理优化思路完整整理一遍核心就是两块分片发送和失败重试机制。无论你是刚接手群发任务的新人还是想把现有发送程序改得更稳的老手这篇文章都值得看完——因为后面写的每一个坑我都亲眼见过别人踩进去。1. 别再for循环调接口一次群发被打到限流的经历先讲个真实场景。某运营团队要群发3万条服务通知负责开发的同事拿到OpenID列表第一反应就是用一个for循环一条一条去调群发消息API接口。前2000条跑得还算正常越往后单次请求耗时越长到第6000条左右接口开始返回限流错误码再往后干脆大面积超时。最后统计下来真正发送成功的不到一半而且因为代码里完全没有记录哪些成功哪些失败连补发的依据都没有只能让用户自己反馈“是不是漏了消息”。这个case看着很初级但很多人第一次做群发就是这么干的。直接循环调接口的问题有三个浪费了最宝贵的资源——请求次数。群发接口普遍有单次调用上限比如一次最多传1000个用户ID。你一条一条发相当于把一个请求能干完的活拆成1000个请求白白消耗调用配额。高频调用必然触发限流。接口侧通常有每分钟或每秒钟的调用次数限制for循环这种“短平快”的调用方式非常容易瞬间把QPS顶上去然后被限流拦截反而拖慢整体进度。没有任何失败补偿。哪条失败了、为什么失败、要不要重试、重试会不会重复发送全都没有设计。一旦出错数据就对不上。所以群发消息API接口的批量处理不能靠蛮力得靠策略。1.1 群发接口普遍存在的三条限制做优化之前先摸清接口的天花板。群发类的API接口通常有三道限制不管具体业务怎么设计大体上都跑不出这几个维度限制维度典型表现影响单次调用数量上限一次最多传N个用户ID超出报参数错误决定你每条请求能“打包”多少人单位时间调用次数限制每秒/每分钟只能调用M次超出报限流错误决定你的并发度上限单用户接收频控同一个用户一天最多接收K条群发影响数据清洗和去重逻辑你接的具体接口文档里一定写明了这些数值。我的建议是开发前先把这三个参数确认清楚做成配置项而不是硬编码在代码里。后面做分片、做并发控制、做重试策略全都要用到这三个数。2. 分片策略设计批大小和并发度不是拍脑袋定的分片发送的核心思路很简单把几十万用户ID拆成若干个“批”每一批作为一次API调用的参数然后用不超过接口限制的并发度去发。但“拆多少一批、同时跑几个线程”这两个问题很多项目是拍脑袋定的结果不是太慢就是被限流。2.1 单批数量受什么约束单批数量首先要看接口单次调用上限。假设上限是500你最好用400左右作为实际批大小留出余量。为什么留余量因为你的服务可能还有其他业务在调同一个接口账号维度的频控是共享的顶着上限发风险太高。另一个约束是单请求的处理耗时。同一批用户数量越多接口侧处理时间越长连接超时的概率越大。一批500个用户单请求耗时可能1.2秒一批100个可能只要300毫秒。你不能只看数量不看延迟。实际项目中我一般这样定批大小batchSize 接口单次上限 × 0.8同时保证单请求耗时压测下来P99不超过2秒两个条件取交集。比如上限1000压测发现800一批要1.8秒OK但如果发现500一批就要2.5秒那就得往下降到400。2.2 并发度计算用耗时和QPS反推并发度不是越大越好。接口QPS上限是5你开50个线程并发结果就是大量请求排队等待限流放行白白把资源耗在阻塞上。正确的做法是用一个简单的关系式去推并发度 ≈ QPS上限 × 单请求平均耗时这是典型的Little法则系统能维持的在途请求数等于吞吐率乘以每个请求的驻留时间。举个例子接口QPS上限是5单请求平均耗时1.2秒那么并发度就是5×1.2≈6。意思是理想状态下6个并发线程同时在飞刚好既能跑满QPS又不会造成严重排队。考虑到网络波动和突发情况我一般在计算值上再乘0.6到0.8的安全系数。比如算出来6实际用4到5。如果你用的是线程池这4到5就是这个线程池的核心线程数。生产环境的参数要以压测为准但计算值作为起点是完全够用的。2.3 一个够用的分片工具类分片逻辑不复杂注意一个坑就行subList返回的是原列表的视图不是副本。如果你在后续流程里还要修改或并发读取原列表视图会出问题。所以分片时最好用new ArrayList(...)包一层生成独立副本。public class BatchPartitionerT { private final int batchSize; public BatchPartitioner(int batchSize) { this.batchSize batchSize; } public ListListT partition(ListT source) { if (source null || source.isEmpty()) { return Collections.emptyList(); } int batchCount (source.size() batchSize - 1) / batchSize; ListListT result new ArrayList(batchCount); for (int i 0; i source.size(); i batchSize) { int end Math.min(i batchSize, source.size()); // 用 new ArrayList 包裹避免 subList 的视图问题 result.add(new ArrayList(source.subList(i, end))); } return result; } }用起来就是ListListString batches new BatchPartitioner(400).partition(openIdList);到这里只是把列表切好真正的重头戏是切片之后怎么发、失败了怎么补救。3. 失败重试机制从“重试三次”到带退避的重试队列很多人写完分片发送会在代码里加一个简单重试for (int i 0; i 3; i) { try { send(batch); break; } catch (Exception e) { // ignore } }这种“重试三次”的写法在真实环境里反而会闯祸。所有失败都用同样的节奏重试接口正在限流时你偏要密集重试等于火上浇油。真正可靠的重试机制要从三个维度设计区分失败类型、设定退避策略、保证幂等。3.1 先分清哪些失败值得重试收到接口返回结果后第一件事不是重试而是判断这个失败有没有重试的价值。我的做法是把失败原因分成三类失败类型典型表现处理策略可重试的临时失败网络超时、连接拒绝、服务端5xx、接口限流按退避策略重试永久业务失败参数错误、无权限、账号被禁用、业务校验不过立即放弃记录原因结果不确定请求发出但响应超时不知道服务端是否处理成功按幂等策略重试或进入人工确认名单判断逻辑不要散落在业务代码里最好统一封装成一个错误分类器。比如限流错误码属于可重试但要走更长的退避参数错误属于永久失败重试一万次也没用。3.2 退避策略的参数设计与抖动退避策略解决的核心问题是失败之后等多久再重试重试之间隔多久。固定间隔只适合极短暂的抖动更好的方案是指数退避加抖动。指数退避的公式是delay base × 2^attempt假设base是2秒第0次失败后等2秒第1次重试失败后等4秒第2次后等8秒。这样既不会在限流期间疯狂试探也不会等得太久让用户着急。抖动为什么要加如果没有抖动一批几十个失败批次都按同样的节奏重试会产生“惊群效应”——某个时间点大家同时苏醒再次把接口打爆。抖动就是在计算值基础上加一个随机偏移long delay base * (1L attempt) ThreadLocalRandom.current().nextLong(0, 1000);重试次数不能无限要有一个总预算。比如base2秒最大重试5次那么整个重试过程的最长等待约248163262秒。超过这个预算还没成功就进入死信队列交给人工处理。这些参数同样做成配置参数建议值说明base 初始延迟2秒第一等待基础值multiplier 倍数2指数增长因子maxAttempt 最大重试次数5超过后进死信队列jitter 抖动范围0~1000毫秒避免惊群3.3 幂等设计重试最怕的事是重复发送重试机制最大的隐患不是“重试了还是失败”而是“其实第一次成功了但因为响应超时被判定失败重试导致用户收到两条一模一样的消息”。解决这个问题只有一个可靠手段幂等。每个批次在发起时生成一个全局唯一的消息标识比如batchId或者更细粒度的messageId接口侧根据这个标识判断是否已经处理过如果处理过就直接返回成功结果。如果你的群发接口不支持幂等键那至少要做到把“结果不确定”的失败单独剥离出来不要自动重试进入一个人工确认名单。我见过太多线上事故起因都是开发默认“重试安全”结果把接口的不确定性忽略了。另外批次状态落库也是幂等的重要保障。发送前把批次状态置为“发送中”成功后才改成“成功”。重试调度器在发起下一轮重试前先检查数据库里的状态如果已经是成功就不需要再发。这个检查必须在事务保护下做否则并发重试会同时放行。3.4 重试调度器的代码骨架用ScheduledExecutorService做定时重试调度是最轻量的方案不需要引入额外的消息队列组件。核心思路是发送失败后把任务包装成RetryTask根据退避策略算出延迟时间重新提交到调度器。public class RetryDispatcher { private final ScheduledExecutorService scheduler Executors.newScheduledThreadPool(2); private final RetryPolicy policy; private final BatchSender sender; public RetryDispatcher(RetryPolicy policy, BatchSender sender) { this.policy policy; this.sender sender; } public void submit(RetryTask task) { SendResult result sender.sendBatch(task.getBatch()); if (result.isSuccess()) { task.markSuccess(); return; } if (!result.isRetryable() || task.getAttempt() policy.getMaxAttempt()) { task.markDead(终态失败不再重试); return; } task.incrementAttempt(); long delay policy.nextDelay(task.getAttempt()); scheduler.schedule(() - submit(task), delay, TimeUnit.MILLISECONDS); } }这个骨架把三类失败都照顾到了成功就结束可重试就按策略退避不可重试直接进死信。4. 工程落地线程池、限流器与批次状态管理策略设计好了接下来是落地。这一节讲实际代码里最容易被写歪的三个部分线程池参数、限流实现、状态追踪。4.1 线程池参数怎么填分片发送用ThreadPoolExecutor参数不能抄网上的默认值。前面算出的并发度就是核心线程数最大线程数最好和核心线程数一致——群发任务是IO密集型的瓶颈在接口侧不在本地CPU把最大线程数拉大只会提高排队阻塞的概率。队列容量则要结合批次总数估算。int concurrence computeConcurrence(); // 用第2节的公式 int maxPending batchCount * 2; // 最多允许积压批次数的2倍 ThreadPoolExecutor executor new ThreadPoolExecutor( concurrence, concurrence, 0L, TimeUnit.MILLISECONDS, new ArrayBlockingQueue(maxPending), new ThreadPoolExecutor.CallerRunsPolicy());拒绝策略用CallerRunsPolicy是故意为之队列满的时候提交任务的线程自己执行发送逻辑相当于降速但不丢任务。群发任务的实时性要求没那么高任务丢一件就是事故宁可让提交线程慢一点也不能把任务丢掉。4.2 本地限流器还是分布式限流器接口有QPS限制本地要做限流保护。最简单的方案是Guava的RateLimiter创建时设置一个预热周期让速率从初始值平滑提升到目标值避免服务刚启动就把接口打满。RateLimiter limiter RateLimiter.create(5, 1, TimeUnit.SECONDS);这里有个重要前提本地限流只对单实例有效。如果你的群发服务部署了多个实例每个实例各有一个RateLimiter5个实例叠加起来就是25 QPS照样把接口打爆。多实例环境必须用分布式限流器最简单的实现是Redis Lua脚本的令牌桶算法。核心思路是用一个key存储当前令牌数和时间戳每次调用前执行一条Lua脚本原子性地扣减令牌令牌不够就返回限流。这段逻辑写成伪代码大概是这样-- 伪代码示意令牌桶的核心判断 local tokens tonumber(redis.call(GET, key) or 0) local capacity tonumber(ARGV[1]) local refill tonumber(ARGV[2]) -- 每秒补充速率 if tokens 1 then return 0 end redis.call(DECRBY, key, 1) return 1能上分布式限流就别只在本地做这是我从线上事故里换来的教训。4.3 批次状态表与兜底轮询重试调度器再怎么可靠也架不住进程重启、机器宕机这种极端情况。内存里的队列一丢所有未完成任务就全没了。所以批次状态必须落库至少要有这样一张表CREATE TABLE batch_send_task ( id BIGINT PRIMARY KEY AUTO_INCREMENT, batch_no VARCHAR(64) NOT NULL COMMENT 批次编号, status TINYINT NOT NULL COMMENT 0初始 1发送中 2成功 3部分成功 4失败 5死信, total_count INT NOT NULL, success_count INT DEFAULT 0, fail_count INT DEFAULT 0, retry_count INT DEFAULT 0, next_retry_time DATETIME DEFAULT NULL, message_id VARCHAR(128) DEFAULT NULL COMMENT 幂等标识, create_time DATETIME NOT NULL, finish_time DATETIME DEFAULT NULL, KEY idx_status_time (status, next_retry_time) );这张表有两个作用。一是给重试调度器提供状态依据重试前先查状态如果已经是成功就不重复发。二是兜底就算内存队列全丢了后台起一个定时任务扫描status1且next_retry_time小于当前时间的批次重新加载并提交发送。网上说的“至少最后一道防线”指的就是这个兜底轮询。4.4 完整调用链路从收到任务到发送完成把上面的模块串起来一个群发任务的完整生命周期是接收原始用户ID列表先做数据清洗按用户维度去重、剔除拉黑名单和退订用户用BatchPartitioner按配置的批大小分片每个批次插入batch_send_task表状态置为初始值提交到线程池执行发送发送成功更新状态为成功记录成功数和耗时发送失败走RetryDispatcher根据错误类型分类可重试的按退避策略重调度超过最大重试次数状态置为死信生成人工处理清单。主流程的骨架代码大致这样public void scheduleSend(ListString openIds) { ListString cleanedIds dataCleaner.clean(openIds); ListListString batches partitioner.partition(cleanedIds); for (ListString batch : batches) { BatchTask task taskRepository.create(batch); executor.execute(() - retryDispatcher.submit(new RetryTask(task))); } }到这里一套基础可用的群发批处理框架已经成型。但框架能跑通只是第一步线上环境里真正考验人的是各种意外。5. 真实环境里踩过的坑与排查过程最后分享几个我在真实环境里遇到过的问题。不是理论推演都是真金白银换来的经验踩过的人应该能会心一笑。5.1 限流错误码被当成业务失败数据悄悄丢了第一个坑就很隐蔽。某次任务跑完总成功数比预期少了一大截差出来的数据没有进失败队列就这么消失了。排查思路是顺着发送日志往回翻发现大量请求返回了限流错误码但代码里的错误处理分支把它当成“业务拒绝”直接丢弃了。限流明明是临时状态正确的处理是走指数退避重试结果被当成永久失败处置。这个坑的教训有两条错误码分类必须在重试机制之前做扎实用测试用例把每个错误码对应的策略固定下来每次发送任务结束要有一个对账步骤把成功数死信数待重试数和任务总数对齐对不上的立即报警。5.2 超时导致重试用户收到了两条消息典型场景线程池并发8个批次某个批次的请求发到一半网络抖动导致客户端读超时。其实服务端已经把消息处理完了但客户端拿不到响应于是把这个批次标记为失败进入重试队列。等网络恢复重试再次发起那一批用户就收到了两条一模一样的消息。排查时先看了用户投诉的时间段再对照日志里重试记录发现重复消息全部集中在某个“超时批次”。后来给接口加上了幂等键能力每个批次带一个全局唯一的messageId服务端去重。如果你的接口没有幂等能力至少要做到超时和业务失败分开记录超时的批次默认进入人工确认列表不要自动重试。5.3 一次性加载全部用户导致内存溢出50万用户ID一次性从数据库读出来全放到一个List里再一遍遍分片。结果任务跑到一半JVM堆内存告急直接OOM。这属于分片做得对、查询方式没分片。排查后把读取方式改成流式处理用游标逐批从数据库读取每读出一批ID就立即清洗、分片、下发处理完一批释放一批。实在需要全量加载的场景至少要做一个数量上限保护超过上限就强制分批加载。5.4 连接池和日志被高并发拖垮并发度提上去之后新的瓶颈冒出来了数据库连接池默认10个连接多线程批量更新批次状态时连接池被打满线程互相等待整体性能反而下降。调大连接池并且把逐条更新改成批量UPDATE状态更新的压力才压下来。还有一个不起眼但很实际的坑日志。如果每个批次里都循环打印每个用户的发送结果几十万条日志瞬间就能把磁盘写满。排查时发现磁盘占用率飙到95%才意识到是日志太“豪放”。后来改成只打印批次维度的汇总日志单条用户日志只在失败时记录整个系统的日志量降了一个数量级。每个群发需求上线前我固定会做三件事把错误码分类和重试策略再过一遍压测确认并发参数小流量灰度观察限流和重试水位。自己踩过这些坑之后我越来越觉得这个流程不能省——群发代码写出来不难难的是让它在各种意外条件下还能保持数据准确、不重复、不遗漏。