微信机器人API接口的Java高并发设计:从连接池到消息异步处理
微信机器人在高频场景(如群聊监控、自动回复)下需支撑数千QPS,若采用同步HTTP调用与阻塞IO,极易因线程耗尽或连接瓶颈导致服务雪崩。本文基于wlkankan.cn.bot包,从HTTP连接池复用、异步非阻塞调用到消息队列削峰,构建端到端高并发架构。
高性能HTTP客户端:复用连接池
使用Apache HttpClient 5.x配置多路复用连接池:
package wlkankan.cn.bot.client;
import org.apache.hc.client5.http.async.HttpAsyncClients;
import org.apache.hc.client5.http.config.ConnectionConfig;
import org.apache.hc.client5.http.impl.async.CloseableHttpAsyncClient;
import org.apache.hc.client5.http.impl.nio.PoolingAsyncClientConnectionManager;
import org.apache.hc.core5.util.TimeValue;
import java.util.concurrent.TimeUnit;
public class WeComHttpClient {
private static final CloseableHttpAsyncClient CLIENT;
static {
ConnectionConfig connConfig = ConnectionConfig.custom()
.setConnectTimeout(TimeValue.of(3, TimeUnit.SECONDS))
.setSocketTimeout(TimeValue.of(10, TimeUnit.SECONDS))
.build();
PoolingAsyncClientConnectionManager cm = new PoolingAsyncClientConnectionManager();
cm.setMaxTotal(200); // 总连接数
cm.setDefaultMaxPerRoute(50); // 每个host最大连接
CLIENT = HttpAsyncClients.custom()
.setConnectionManager(cm)
.setDefaultRequestConfig(org.apache.hc.client5.http.config.RequestConfig.custom()
.setResponseTimeout(TimeValue.of(10, TimeUnit.SECONDS))
.build())
.build();
CLIENT.start();
}
public static CloseableHttpAsyncClient getInstance() {
return CLIENT;
}
}

异步发送消息:CompletableFuture封装
package wlkankan.cn.bot.service;
import wlkankan.cn.bot.client.WeComHttpClient;
import org.apache.hc.core5.http.ContentType;
import org.apache.hc.core5.http.io.entity.StringEntity;
import org.apache.hc.core5.http.message.BasicClassicHttpRequest;
import java.util.concurrent.CompletableFuture;
public class AsyncMessageSender {
public CompletableFuture<String> sendText(String accessToken, String userId, String content) {
String url = "https://qyapi.weixin.qq.com/cgi-bin/message/send?access_token=" + accessToken;
String json = String.format("""
{"touser":"%s","msgtype":"text","text":{"content":"%s"}}
""", userId, content);
BasicClassicHttpRequest request = new BasicClassicHttpRequest("POST", url);
request.setEntity(new StringEntity(json, ContentType.APPLICATION_JSON));
CompletableFuture<String> future = new CompletableFuture<>();
WeComHttpClient.getInstance().execute(request, new FutureCallback<>() {
@Override
public void completed(org.apache.hc.core5.http.ClassicHttpResponse response) {
try {
String body = org.apache.hc.core5.util.EntityUtils.toString(response.getEntity());
future.complete(body);
} catch (Exception e) {
future.completeExceptionally(e);
}
}
@Override
public void failed(Exception ex) {
future.completeExceptionally(ex);
}
@Override
public void cancelled() {
future.cancel(false);
}
});
return future;
}
}
消息队列削峰:Disruptor高性能环形缓冲
引入LMAX Disruptor避免锁竞争:
package wlkankan.cn.bot.queue;
import com.lmax.disruptor.EventHandler;
import com.lmax.disruptor.RingBuffer;
import com.lmax.disruptor.dsl.Disruptor;
import wlkankan.cn.bot.model.SendMessageEvent;
import java.util.concurrent.Executors;
public class MessageQueue {
private final Disruptor<SendMessageEvent> disruptor;
private final RingBuffer<SendMessageEvent> ringBuffer;
public MessageQueue(EventHandler<SendMessageEvent> handler) {
this.disruptor = new Disruptor<>(
SendMessageEvent::new,
1024 * 64, // 环大小,必须为2的幂
Executors.defaultThreadFactory(),
com.lmax.disruptor.ProducerType.MULTI,
new com.lmax.disruptor.BlockingWaitStrategy()
);
disruptor.handleEventsWith(handler);
disruptor.start();
this.ringBuffer = disruptor.getRingBuffer();
}
public void publish(String accessToken, String userId, String content) {
long sequence = ringBuffer.next();
try {
SendMessageEvent event = ringBuffer.get(sequence);
event.reset(accessToken, userId, content);
} finally {
ringBuffer.publish(sequence);
}
}
public void shutdown() {
disruptor.shutdown();
}
}
package wlkankan.cn.bot.model;
public class SendMessageEvent {
private String accessToken;
private String userId;
private String content;
public void reset(String accessToken, String userId, String content) {
this.accessToken = accessToken;
this.userId = userId;
this.content = content;
}
// getters
public String getAccessToken() { return accessToken; }
public String getUserId() { return userId; }
public String getContent() { return content; }
}
消费者:批量处理与背压控制
package wlkankan.cn.bot.consumer;
import wlkankan.cn.bot.model.SendMessageEvent;
import wlkankan.cn.bot.service.AsyncMessageSender;
import com.lmax.disruptor.EventHandler;
public class MessageConsumer implements EventHandler<SendMessageEvent> {
private final AsyncMessageSender sender = new AsyncMessageSender();
@Override
public void onEvent(SendMessageEvent event, long sequence, boolean endOfBatch) {
sender.sendText(event.getAccessToken(), event.getUserId(), event.getContent())
.exceptionally(ex -> {
// 记录失败日志,可重试或告警
System.err.println("Send failed: " + ex.getMessage());
return null;
});
}
}
控制器集成
package wlkankan.cn.bot.controller;
import wlkankan.cn.bot.queue.MessageQueue;
import org.springframework.web.bind.annotation.*;
@RestController
@RequestMapping("/bot")
public class BotController {
private final MessageQueue queue = new MessageQueue(new wlkankan.cn.bot.consumer.MessageConsumer());
@PostMapping("/send")
public String sendMessage(@RequestParam String token,
@RequestParam String user,
@RequestBody String content) {
queue.publish(token, user, content);
return "queued";
}
}
通过wlkankan.cn.bot模块实现的连接池复用、异步非阻塞调用与Disruptor无锁队列,系统在万级QPS下仍保持低延迟与高吞吐,有效支撑企业级微信机器人场景。




