凌晨两点的告警:消息积压 12 万条
6 月 18 号大促那晚,我睡到两点被电话叫醒。监控群里刷的是一条 RocketMQ 的告警:order_pay_topic 消费 TPS 从平时的 800 掉到了 30,积压量 12 万条还在涨。
赶紧登机器看日志,消费者进程活着,CPU 只有 11%,但日志里满屏都是这个:
2019-06-18 02:07:43.118 [ConsumeMessageThread_9] ERROR c.x.mq.OrderPayConsumer - 处理消息失败
java.util.concurrent.RejectedExecutionException: Task com.xxx.mq.OrderPayConsumer$$Lambda$412/0x00000007c0a1e040
rejected from java.util.concurrent.ThreadPoolExecutor@5f184fc6
[Running, pool size = 20, active threads = 20, queued tasks = 2000, completed tasks = 1837442]
at java.util.concurrent.ThreadPoolExecutor$AbortPolicy.rejectedExecution(ThreadPoolExecutor.java:2063)
at java.util.concurrent.ThreadPoolExecutor.reject(ThreadPoolExecutor.java:830)
at java.util.concurrent.ThreadPoolExecutor.execute(ThreadPoolExecutor.java:1379)
at com.xxx.mq.OrderPayConsumer.consumeMessage(OrderPayConsumer.java:57)
任务被拒了。看这行 [Running, pool size = 20, active threads = 20, queued tasks = 2000],池子满了 20 个线程全在忙,队列 2000 也塞满了,第 2001 个任务进来就被扔掉。而 RocketMQ 客户端收到异常后会返回 RECONSUME_LATER,消息重试,重试又失败,于是越积越多。
更要命的是我这行代码:
@Component
@RocketMQMessageListener(topic = "order_pay_topic", consumerGroup = "order_pay_cg")
public class OrderPayConsumer implements RocketMQListener<String> {
private final ExecutorService bizPool = Executors.newFixedThreadPool(20);
@Override
public void onMessage(String body) {
// 丢进自己的业务线程池就返回,让 RocketMQ 的线程继续拉消息
bizPool.execute(() -> handlePaySuccess(body));
}
}
我当初的打算是"消费线程只做转发,业务处理在另一个池子里跑,提高吞吐"。结果业务处理要调优惠券和积分两个下游,平均耗时 240 毫秒,20 个线程的极限吞吐是 20 / 0.24 ≈ 83 TPS。而大促期间生产端峰值是 900 TPS。缺口摆在那儿,队列只是延缓了爆发时间。
先搞清楚拒绝了会怎样
JDK 8 的 ThreadPoolExecutor 内置了四种拒绝策略,都在它的内部类里,行为差别很大:
| 策略 | 行为 | 会丢任务吗 |
|---|---|---|
| AbortPolicy(默认) | 抛 RejectedExecutionException | 丢,但调用方能感知 |
| CallerRunsPolicy | 让提交任务的线程自己跑 | 不丢 |
| DiscardPolicy | 静默丢弃,什么都不做 | 丢,且无声无息 |
| DiscardOldestPolicy | 丢掉队列头那个,再重试提交 | 丢,丢的是最老的 |
源码就几行,看一遍就记住了:
public static class AbortPolicy implements RejectedExecutionHandler {
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
throw new RejectedExecutionException("Task " + r.toString() +
" rejected from " + e.toString());
}
}
public static class CallerRunsPolicy implements RejectedExecutionHandler {
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
if (!e.isShutdown()) {
r.run(); // 注意:直接 run(),不是另起线程
}
}
}
public static class DiscardOldestPolicy implements RejectedExecutionHandler {
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
if (!e.isShutdown()) {
e.getQueue().poll(); // 扔掉队头
e.execute(r); // 再试一次,可能又被拒
}
}
}
默认的 AbortPolicy 其实还算"仁慈",至少它喊了一声。DiscardPolicy 的 rejectedExecution 方法体是空的,任务凭空消失,日志里一点痕迹都没有。我们组另一个同事在异步写埋点的池子上用了它,埋点数据少了 30% 半个月才被发现。
Executors.newFixedThreadPool(20) 这个写法本身就埋了雷,它内部是:
public static ExecutorService newFixedThreadPool(int nThreads) {
return new ThreadPoolExecutor(nThreads, nThreads,
0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<Runnable>()); // 无界队列!
}
等等,无界队列怎么会触发拒绝?因为我这不是无界队列,实际代码里为了"防止内存爆掉"手动换成了 ArrayBlockingQueue(2000)。但 newFixedThreadPool 默认那个 LinkedBlockingQueue 容量是 Integer.MAX_VALUE,任务会一直堆到 OOM 为止,那更可怕。
CallerRunsPolicy 怎么起到反压作用
问题的本质是:上游 900 TPS 往里灌,下游只能吃 83 TPS,多出来的必须有个去处。要么排队(队列会涨),要么丢弃(会丢数据),要么让上游慢下来。第三种才是正解,而 CallerRunsPolicy 天然就是干这个的。
它的逻辑是:池子满了,提交任务的那个线程(这里是 RocketMQ 的 ConsumeMessageThread)自己把任务跑完再返回。这一跑就是 240 毫秒,期间它不会去拉新消息,Broker 那边的消费进度也就不推进。等它忙完,才继续拉下一条。等效于把消费速率压到下游能承受的水平。
我改完之后的版本:
private final ThreadPoolExecutor bizPool = new ThreadPoolExecutor(
20, 20, 0L, TimeUnit.MILLISECONDS,
new ArrayBlockingQueue<>(2000),
new ThreadFactoryBuilder().setNameFormat("biz-pay-%d").build(),
new ThreadPoolExecutor.CallerRunsPolicy());
这里有个细节:CallerRunsPolicy 生效的前提是提交者和执行者是不同的线程。如果提交方是 Netty 的 IO 线程或者 Tomcat 的 acceptor 线程,让它们去跑业务逻辑会把整个连接层拖死。RocketMQ 的消费线程本来就专职消费,让它顶一会儿没问题。
改完压测了一把,用 900 TPS 灌 5 分钟:
| 策略 | 积压峰值 | 任务丢失 | 消费端 CPU | 5 分钟总吞吐 |
|---|---|---|---|---|
| AbortPolicy | 12 万+ | 8734 条(重试后仍失败) | 11% | 30 TPS |
| CallerRunsPolicy | 2100(队列上限) | 0 | 63% | 79 TPS |
确实不丢消息了,但吞吐只有 79 TPS,积压消不掉。这只是止血,根本办法还是提并发。
自定义策略:落盘 + 补偿重试
后来我把线程数提到 60(下游压测能扛 260 TPS),同时写了个自定义策略兜底。思路是:能反压就反压,实在压不住的落盘到本地文件,由一个单独的定时任务慢慢回放。
public class DumpAndRetryPolicy implements RejectedExecutionHandler {
private static final Logger log = LoggerFactory.getLogger(DumpAndRetryPolicy.class);
private final String dumpDir;
public DumpAndRetryPolicy(String dumpDir) {
this.dumpDir = dumpDir;
}
@Override
public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) {
// 提交方是业务线程,先尝试反压,避免落盘 IO 阻塞调用链
if (Thread.currentThread().getName().startsWith("biz-")) {
if (!executor.isShutdown()) {
r.run();
return;
}
}
// 提交方已经是池内线程(比如嵌套提交),不能再 run,否则递归爆栈
log.warn("pool saturated, dump task. pool={}, queue={}",
executor.getPoolSize(), executor.getQueue().size());
dumpToFile(r, executor);
}
private void dumpToFile(Runnable r, ThreadPoolExecutor executor) {
File dir = new File(dumpDir);
if (!dir.exists() && !dir.mkdirs()) {
log.error("create dump dir failed: {}", dumpDir);
return;
}
File file = new File(dir, "rejected-" + LocalDate.now() + ".json");
try (FileWriter fw = new FileWriter(file, true);
PrintWriter pw = new PrintWriter(fw)) {
pw.println(TaskSerializer.toJson(r, executor));
} catch (IOException e) {
log.error("dump task failed", e);
}
}
}
那个"判断线程名"的分支看着别扭,但踩过坑才写得出来。我们第一版没这个判断,结果 handlePaySuccess 内部又往同一个池子提交了子任务,池满时子任务被拒,r.run() 又在当前线程(池内线程)执行,子任务里再提交……形成递归,最后 StackOverflowError。判断当前线程是不是池内的,是就直接落盘。
回放的定时任务用 Spring 的 @Scheduled,每 30 秒读一次文件,限速 50 TPS 慢慢补:
@Scheduled(fixedDelay = 30_000)
public void replay() {
File file = new File(dumpDir, "rejected-" + LocalDate.now() + ".json");
if (!file.exists()) {
return;
}
// 先 rename 成 .processing,防止和正在写的文件冲突
File processing = new File(file.getAbsolutePath() + ".processing");
if (!file.renameTo(processing)) {
return;
}
int success = 0, fail = 0;
try (BufferedReader br = new BufferedReader(new FileReader(processing))) {
String line;
while ((line = br.readLine()) != null) {
if (rateLimiter.tryAcquire()) { // Guava RateLimiter,50/s
try {
handlePaySuccess(line);
success++;
} catch (Exception e) {
fail++;
log.error("replay failed: {}", line, e);
}
}
}
} catch (IOException e) {
log.error("read dump file failed", e);
}
log.info("replay done, success={}, fail={}", success, fail);
processing.delete();
}
大促后统计,整个过程落盘了 412 条任务,回放全部成功,没有一条消息丢。
小结
Executors.newFixedThreadPool用的是无界LinkedBlockingQueue,生产环境自己 newThreadPoolExecutor,队列长度一定要显式给。- 默认
AbortPolicy会抛异常,在消息消费场景会触发重试风暴;DiscardPolicy静默丢任务,最危险,别用。 CallerRunsPolicy是最好的默认选择,它把压力推回上游形成反压。前提:提交线程不能是 IO 线程。- 自定义策略里如果打算调用
r.run(),一定要判断当前线程是不是池内线程,否则可能递归爆栈。 - 反压只能保不丢,消不了积压。真要解决还是得扩并发或者优化单次耗时——我们后来把两个下游调用改成并行,单次耗时从 240 毫秒降到 130 毫秒,吞吐直接翻倍。