Administrator
发布于 2019-06-21 / 2890 阅读
43

线程池的拒绝策略怎么选?一次消息队列积压的教训

凌晨两点的告警:消息积压 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 其实还算"仁慈",至少它喊了一声。DiscardPolicyrejectedExecution 方法体是空的,任务凭空消失,日志里一点痕迹都没有。我们组另一个同事在异步写埋点的池子上用了它,埋点数据少了 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 分钟:

策略积压峰值任务丢失消费端 CPU5 分钟总吞吐
AbortPolicy12 万+8734 条(重试后仍失败)11%30 TPS
CallerRunsPolicy2100(队列上限)063%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,生产环境自己 new ThreadPoolExecutor,队列长度一定要显式给。
  • 默认 AbortPolicy 会抛异常,在消息消费场景会触发重试风暴;DiscardPolicy 静默丢任务,最危险,别用。
  • CallerRunsPolicy 是最好的默认选择,它把压力推回上游形成反压。前提:提交线程不能是 IO 线程。
  • 自定义策略里如果打算调用 r.run(),一定要判断当前线程是不是池内线程,否则可能递归爆栈。
  • 反压只能保不丢,消不了积压。真要解决还是得扩并发或者优化单次耗时——我们后来把两个下游调用改成并行,单次耗时从 240 毫秒降到 130 毫秒,吞吐直接翻倍。

参考