尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

从微信群发接口实战看Java后端批量处理与性能调优

从微信群发接口实战看Java后端批量处理与性能调优 今年年初调完一个微信群发消息的接口项目趁记忆还热把Java后端在批量数据处理和性能调优上踩过的坑、试对的路子完整复盘一遍。先交代一下背景业务方要通过微信公众号服务号的模板消息把几十万用户按标签、活跃时间和地域分成不同批次做精准触达。单次任务用户量在20万到50万要求2小时内发完并且要做到失败可重试、不重复发送。这个量级对数据库、Java服务以及微信接口通道来说都不是“写个循环”就能糊弄过去的。这篇文章就围绕这类“群发消息API接口”的Java后端实现把批量数据处理和性能调优的思路一点点拆开从任务怎么拆、数据怎么读、线程池怎么配到微信接口限流、幂等重试、线上压测和问题排查每一段都会给出可以直接参考的方案和代码。适合正在做消息推送、任务调度、批量接口调用的后端同学哪怕你没有接触过微信生态里面的批量处理思路放到其他第三方API对接场景也一样适用。1. 项目背景微信群发接口背后的批量难题1.1 一个真实的群发任务长什么样上线前的需求听上去很简单运营同学建好用户分组点击“群发”系统把模板消息发给组内所有用户。但落到Java后端这个需求会变成一条很长的数据链路——从数据库里查出几十万目标用户把每个用户的姓名、订单号、优惠券信息填进模板再逐个调用微信API推送最后把每条结果写回库。我当时负责的模块是“任务编排中枢”。每次群发任务生成后系统要处理这么几件事从用户表按标签/条件圈定目标人群生成任务明细。根据微信接口的能力把明细拆成多个批次。通过线程池并发调用微信模板消息接口。实时记录每一条消息的发送状态支持失败后重新推送。看起来每个环节都不复杂但一旦数据量大起来问题就全冒出来了。最痛苦的不是“调接口”而是“怎么让大批量数据在有限时间内处理完、同时不拖垮数据库和内存”。1.2 Java后端批量处理的三个敌人做这类项目我最大的感受是性能瓶颈往往不在业务代码有多复杂而在三个“看不见的敌人”。第一个是内存。假如一次导入50万用户每个用户封装成对象带着openid、昵称、参数Map粗算下来一个对象可能要占几百字节。如果全部加载进一个List光这一波就是几百MB再来几个并发任务JVM堆直接告急GC频繁到CPU飙高。很多人写批量任务习惯“先查全量再慢慢处理”这是最容易OOM的写法。第二个是数据库连接。假设你要循环往数据库里写状态每处理完一条消息就update一下50万条就是50万次单条SQL。即便连接池配了50个连接数据库也扛不住这种“小碎步”请求。更别说还要同时处理查询、分页行锁和IO很快会把数据库拖住。第三个是第三方接口的限制。微信官方接口不是你想调就能无限调它有access_token有效期、频率限制、并发限制。很多时候不是我们服务处理得慢而是第三方接口成了“瓶颈阀门”。如果把线程池调得很大去硬怼等着你的就是封禁提醒和大量报错。所以群发消息的Java后端从设计第一天就要把这三个敌人考虑进去否则后面性能调优就是拆东墙补西墙。1.3 为什么性能调优必须从设计阶段开始我见过不少项目前期图快用最简单的方式把功能跑通等到数据量上来才到处“调优”。结果发现数据结构已经定死批量拆分改不动状态字段缺失幂等只能靠Redis锁硬补最后只能推倒重来。微信群发任务这种场景性能调优不该是“出问题再优化”而应该在设计阶段就先问自己几个问题数据量级是多少目标峰值是多少第三方接口的QPS上限是多少数据库能不能承受批量写入任务执行中进程重启了、机器挂了怎么续跑消息有没有可能重复发送业务允许多少误差这些问题想清楚了后面写代码就是顺水推舟。下面我按自己的实际设计步骤从整体数据流开始拆解。2. 整体设计与数据流拆解2.1 任务拆分和状态管理怎么做批量处理项目里我习惯把“一次群发”抽象成三个层面的数据任务、批次、明细。任务表记录一次群发的整体信息比如任务名称、类型、状态待执行、执行中、已完成、部分失败、目标总人数、成功/失败数量、开始结束时间。批次表记录拆出来的每个子批次比如第1批、第2批每批包含一批用户ID。批次表用来支持断点续跑和并发控制。明细表记录每个用户的具体发送结果包含用户ID、openid、消息参数、状态待发送、发送中、成功、失败、重试中、错误码、重试次数。这个结构的好处是每一层都可以独立查询和统计。任务挂了可以根据批次表找到未完成的批次从断点继续跑用户反馈某条没收到直接查明细表定位原因。状态机我一般这样设计任务从“待执行”进入“执行中”所有批次处理完后进入“已完成”。如果存在失败明细且还在重试次数内任务标记为“部分失败”等待补偿任务扫描重发。超过最大重试次数的失败明细保留错误信息供人工排查。这个设计能给后续的“可重试、幂等、对账”打好基础。没有状态管理任何性能调优都是虚的。2.2 每批数据量多大才合适“分批”是批量数据处理的基本手段但“每批拆多大”非常有讲究。拆太大内存吃紧单元处理时间太长一个批次失败要重做很多功拆太小线程切换和SQL查询次数太多效率反而低。我当时是按“内存占用 第三方接口限流 数据库压力”三个维度来定的假设每个用户对象加消息参数约占500B单批次2000人20个线程并发处理驻留内存约2000*500B*20 20MB对JVM来说很安全。单批2000人意味着要调2000次微信模板消息接口。按微信常用限制每秒10次来算2000次要200秒一个批次处理时间过长。所以批大小不能只看内存还要看下游吞吐。最后我选了“以1000人为一个处理单元内部再按并发上限调度”折中下来内存、数据库、接口三方面都能接受。批大小其实没有标准答案只能结合自己的数据模型和接口限制去压测。但一个通用的原则是批大小优先保证内存安全线程池大小优先保证下游接口不被压垮。2.3 异步线程池还是消息队列别一上来就上MQ做群发任务很多人第一反应是上RocketMQ/Kafka。MQ确实能解耦和削峰但没必要场景都套它。一切都是MQ会带来新的问题消息顺序、消费幂等、失败重投、中间件维护成本。如果只是单机单任务的批量推送用线程池加数据库任务表完全够用。我用线程池就实现了异步处理、限速、失败重试而且代码更简单定位问题也更直接。但如果你们系统里已经有一套成熟的MQ且群发任务可能同时跑好几个我会建议用MQ把“任务拆分”和“任务执行”解耦主服务拆好批次把批次ID发给MQ。消费端拿到批次ID从数据库捞明细去执行发送。这样做的优势是天然支持分布式消费、削峰填谷。劣势是需要额外处理消费幂等。我的建议是团队没有MQ基础就别硬上先把任务表线程池玩熟练有MQ就把批次消息丢进去不要直接把50万条明细一条条丢进去。2.4 数据读取优化为什么批量任务要用游标分页而不是offset分页查询目标用户列表时最容易踩的坑就是写这样的SQLSELECT * FROM user_target WHERE task_id ? ORDER BY id LIMIT 1000 OFFSET 50000;offset一深数据库就要先扫到第50000行再往后取翻到后面每一页都越来越慢。更严重的是如果用户在分页过程中数据发生变化还可能出现重复或漏掉数据。我改成游标分页核心就是记住“上一批最后一条的id”下一批带上这个id作为起点SELECT * FROM user_target WHERE task_id ? AND id ? ORDER BY id LIMIT 1000;配合id上的主键索引无论翻到多后面性能都稳定。Java端的逻辑也简单long lastId 0; int batchSize 1000; while (true) { ListUserTarget users userTargetDao.selectByCursor(taskId, lastId, batchSize); if (users.isEmpty()) { break; } lastId users.get(users.size() - 1).getId(); // 把这个批次交给线程池处理 sendTaskExecutor.submit(new SendBatchTask(users)); }这里的核心是“边读边发”不要等把所有用户都读出来再开始处理。游标分页配合“循环内外移出可复用对象”能在源头上控制内存增长这也是批量任务里很关键的优化点。3. 核心实现Java后端批量处理的关键代码3.1 用户数据的流式读取与组装在设计里我不建议一次性把整批用户全部load到内存再组装模板参数。更好的做法是每次从数据库读取一批组装好之后立即提交给线程池然后继续读下一批。这里有一个容易被忽略的小优化在循环外部创建一次对象循环体内部复用必要组件。// 精简示例游标读取 组装消息体 MessageTemplate template messageTemplateService.getById(templateId); long lastId 0L; while (true) { ListUserTarget users userTargetDao.selectByCursor(taskId, lastId, batchSize); if (users.isEmpty()) { break; } lastId users.get(users.size() - 1).getId(); ListSendTask tasks new ArrayList(users.size()); for (UserTarget user : users) { // 组装模板消息填充用户昵称、订单号等个性化内容 SendTask task new SendTask(); task.setOpenid(user.getOpenid()); task.setMsgData(buildMsgData(user, template)); tasks.add(task); } sendTaskExecutor.submit(new SendBatchTask(taskId, tasks)); }组装消息时不要用大量字符串拼接去拼JSON建议用Map或Jackson的ObjectNode构造结构然后一次性writeValueAsString。50万次字符串拼接GC和内存占用都会很感人。3.2 线程池参数配置与动态调整这块是整个批量接口调用的心脏。我当时用的ThreadPoolExecutor做异步发送参数配置如下ThreadPoolExecutor executor new ThreadPoolExecutor( 8, // corePoolSize 32, // maximumPoolSize 60L, TimeUnit.SECONDS, // 空闲线程存活时间 new LinkedBlockingQueue(64), // 有界队列 new ThreadFactoryBuilder().setNameFormat(wechat-send-pool-%d).build(), new CallerRunsPolicy() // 拒绝策略由提交任务的线程自己执行 );这里每个参数都有讲究corePoolSize8这个任务属于IO密集型大量时间在等待微信接口响应线程数过低会让CPU闲置。maximumPoolSize32不是越大越好。微信接口有频率限制线程数太多会把请求打爆。32是我压测后得出的数如果你接入的其他API并发限制更低这个值还要往下调。LinkedBlockingQueue(64)有界队列是必须的防止任务无限堆积导致OOM。CallerRunsPolicy当队列满了由提交任务的线程自己执行这样等于天然做了一层背压不会丢任务。发送线程自己跑会拖慢主线程的读库速度反过来迫使主线程降速形成“自我保护”。IO密集型线程数有一个经验公式线程数 CPU核心数 * 2 * (1 IO等待占比/CPU计算占比)。但实际不用算那么精确微信接口RT大概率在100~300ms按并发度去压测找出“再增加线程数QPS也不涨反而报错频发”的拐点那就是最大线程数。3.3 微信接口调用的限流与access_token刷新线程池解决了并发但不解决“第三方接口限流”问题。微信API对接口调用频率是有明确“配额”的如果所有线程一股脑全速调用很快会触发接口返回异常。所以我通常会再加一层“信号量”或“令牌桶”做客户端限流。用Semaphore做最简单的并发数限制// 限制最多同时10个请求在飞 private final Semaphore wechatSendLimiter new Semaphore(10); public void sendWechat(SendTask task) { try { wechatSendLimiter.acquire(); WechatResponse resp wechatApi.sendTemplateMessage(task); // 处理结果 } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { wechatSendLimiter.release(); } }如果接口是按“每秒请求数”限流的可以用Guava的RateLimiter本质上就是令牌桶RateLimiter rateLimiter RateLimiter.create(5.0); // 每秒最多5个请求 rateLimiter.acquire();还有access_token的问题。微信接口的access_token有效期是7200秒且获取接口有频率限制。批量发送时最怕每个线程都去调“刷新token”白白消耗频率还容易触发锁。我在项目里的做法是全局缓存access_token放到Redis里带过期时间。设置“提前5分钟续期”即有效期还剩5分钟时直接刷新并更新缓存。多实例部署时要加分布式锁防止多个实例同时刷新同一个token。伪代码大致如下public String getAccessToken(String appId, String appSecret) { String cacheKey wechat:access_token: appId; String token redisTemplate.opsForValue().get(cacheKey); if (StringUtils.hasText(token)) { return token; } // 分布式锁防止多实例并发刷新 RLock lock redissonClient.getLock(cacheKey :lock); lock.lock(); try { token redisTemplate.opsForValue().get(cacheKey); if (StringUtils.hasText(token)) { return token; } token doRefreshAccessToken(appId, appSecret); redisTemplate.opsForValue().set(cacheKey, token, Duration.ofSeconds(7000)); return token; } finally { lock.unlock(); } }这个细节很多人会漏掉一旦token过期且没做缓存预热大批量发送时就会看到一大片401或40001错误。3.4 结果回写与失败重试的幂等设计群发消息接口调用完必须把结果回写到数据库。回写同样不能一条条update而是按批次批量更新。更加关键的是“幂等”网络抖动、服务重启、微信接口返回超时都可能导致同一条消息被发送两次。为了避免重复推送我在明细表上加了一个唯一键(task_id, openid)。每次发送前先尝试把明细状态从“待发送”更新为“发送中”执行UPDATE message_detail SET status SENDING, update_time now() WHERE task_id ? AND openid ? AND status PENDING;如果影响行数为0说明这条已经被其他线程“抢”走了就不能再发。发送成功后再更新为“SUCCESS”。失败需要重试时要把状态先改回“PENDING”并把重试次数加一。这样即使同一个任务被重复调度或者多个消费者同时处理同一批用户也不会对同一条消息发两次。这是批量任务里最不能省的一环。3.5 数据库批量写入的JDBC调优前面说过状态回写如果一条条执行SQL性能会很差。要改成批量写入。如果直接用JDBC需要开启rewriteBatchedStatementsMySQL才会把多条insert合并成一条多值SQL。我用的是MyBatis配合批量提交也能达到类似效果// SqlSessionTemplate批量模式 SqlSession sqlSession sqlSessionTemplate.getSqlSessionFactory().openSession(ExecutorType.BATCH); try { MessageDetailMapper mapper sqlSession.getMapper(MessageDetailMapper.class); for (SendTask task : taskList) { mapper.insert(task); } sqlSession.commit(); } finally { sqlSession.close(); }如果是自己拼JDBC连接串一定要加上这个参数jdbc:mysql://localhost:3306/xxx?rewriteBatchedStatementstrueuseServerPrepStmtstrue批量插入的批大小也不是越大越好。我实测过单批500条到5000条之间性能都比较稳定超过1万条反而会因为网络包太大和数据库事务日志膨胀而变慢。优先选1000条左右比较稳。4. 性能调优实战从压测到上线的全过程4.1 先定性能目标再动手压测线上调优最忌讳“边压边调没有目标”。做群发任务时我先把需求换算成了技术指标需求是2小时内发完50万条消息。换算下来平均每秒需要约70条发送成功。考虑失败重试、数据库写入、网络抖动必须留出至少2倍余量。所以我的目标是单实例至少支撑每秒150次微信API调用并且处理完成后数据库写入不积压。有了这个目标压测才有底。如果压测结果只有每秒50次说明线程池或限流配置不合理得继续调如果达到了每秒200次但这个过程中GC频繁、CPU满负载也要警惕。4.2 压测中遇到的三个典型瓶颈第一次压测我用的是Jmeter模拟群发任务的HTTP入口同时在Java端打印关键批次的耗时。结果暴露了三个问题第一个问题是数据库慢SQL。期初用户列表查询用的是offset分页当页数翻到几万之后单条SQL耗时飙升到2秒。这个好解决换成游标分页后SQL耗时稳定在20ms左右。第二个问题是线程池配置过大反而拖慢整体。我把最大线程数调到64结果微信接口频繁报错重试变多整体吞吐反而下降。后来我把最大线程数降到32加上信号量限流整体吞吐反而上来了。这验证了一个道理下游接口能承受多少并发才是真正的上限。第三个问题是消息明细批处理时GC压力大。刚开始用ExecutorType.BATCH但每个批次提交后没有及时清理缓存内存里堆积了大量待写对象老年代一直涨。后来我强制每5000条提交一次并周期性清理上下文GC立刻变得平稳。4.3 调优三板斧批大小、线程数、连接池把压测发现的坑补齐之后我总结出了群发类项目的“调优三板斧”批大小读取批1000写入批5000。读取太大会增加驻留内存写入太大会锁表太久。线程数核心8最大32队列64。如果接口限流更低比如微信允许每秒2次那么最大线程数要压到更低。数据库连接池Druid或HikariCP最小连接数建议设8最大连接数设20。并不是越大越好连接数超过数据库能同时处理的范围后反而增加上下文切换。这里给一个参考表方便不同接口限制下快速调整第三方接口允许的并发量建议最大线程数信号量/限流值58410168303220505040核心原则是线程池最大线程数 下游接口最大并发 少量余量信号量比线程池再紧一点形成两级缓冲。5. 常见问题与排查技巧实录5.1 任务队列满了怎么办使用有界队列后任务高峰期可能出现线程池拒绝任务。如果你用的是CallerRunsPolicy提交任务的线程会被迫自己执行任务看起来是“卡住”其实是在背压保护。有个更稳妥的处理方式准备一个“暂存表”。当线程池队列满时先把多余批次的状态更新为“待执行”不继续提交而是等一批任务跑完后由调度线程重新扫描暂存表继续提。这样既不会丢任务也不会打爆下游。5.2 微信接口报错返回码怎么归类处理微信接口错误码很多不能所有code都无脑重试。如果某个code代表“参数错误”重试多少次都一样。我的处理思路是分成三类可重试系统繁忙-1、token失效40001、频率限制45009。这类等一会儿再发。不可重试参数错误40003、模板不存在40037等。这类直接把明细标记为失败并记录原因。未知先记录原始返回给到告警人工介入。线上群发时我用一个简单的映射表分类错误码配合重试次数上限默认3次间隔指数退避有效避免了无效重试把接口带宽耗尽。5.3 分布式环境下如何避免重复发送如果服务多实例部署线程池是每个实例一套不同实例可能同时扫到同一个任务批次。我做了两层防护数据库唯一索引(task_id, openid)同一个用户只会有一条发送明细。发送前用Redis分布式锁锁住批次ID保证同一批只能被一个实例消费。Redis锁伪代码String lockKey batch:lock: batchId; boolean locked redisTemplate.opsForValue().setIfAbsent(lockKey, 1, Duration.ofMinutes(10)); if (!locked) { // 说明其他实例已在处理跳过 return; } try { processBatch(batchId); } finally { redisTemplate.delete(lockKey); }再加上状态机的那条“PENDING - SENDING”条件更新整个链路就等于上了双重保险。我实测下来分布式场景下基本不会出现重复发送。5.4 发完之后的“对账”技巧群发任务结束后不能只看任务表显示“已完成”就收工。我习惯再加一个对账环节定时任务扫描任务状态对比目标人数、成功数、失败数、重试数是否一致。对失败超过次数上限的明细按用户维度汇总生成“未送达用户清单”给运营。对已发送但微信异步回调显示“用户拒收”的消息再更新一次状态。听起来这一步和性能调优关系不大但对业务来说非常重要。群发不是“调用完接口就结束”消息是否真正送达是运营判断活动效果的关键依据。Java后端在这里要做的是把状态数据闭环起来而不是只盯着发送速率。再分享一点我个人实际操作的体会批量任务最怕的不是“慢”而是“不可恢复”。调优方案做得再漂亮如果任务执行到一半服务重启了没法从断点续跑那才是灾难。我后来养成的习惯是所有批量处理代码先保证“随时可以停、随时可以继续”再谈性能优化。毕竟群发任务这种场景稳定性和数据一致性永远排在吞吐量前面。希望这套从设计到压测、再到线上排查的思路能让你在下次接到类似API批量对接项目时少走几次弯路。
返回列表