Day 1|Java 线程池:从 ThreadPoolExecutor 到 BizeNova AI 异步任务 封面
返回上一级

Day 1|Java 线程池:从 ThreadPoolExecutor 到 BizeNova AI 异步任务

2026.09.02
4
📂 技术
# Java
# 线程池
# BizeNova
# 异步任务

今天的目标不是再背一轮“七大参数、四大拒绝策略”,而是建立一套能落到真实项目里的判断框架:线程池解决的是进程内并发资源治理,不等于可靠任务系统;异步提高的是请求响应速度,不会自动带来任务不丢、结果不重、故障可恢复。

本文先从 ThreadPoolExecutor 的执行模型讲起,再拆解 BizeNova 当前的企业 AI 评估链路,最后给出不急着上 MQ 的可靠化方案,以及可以直接用于面试的问答。

一、为什么需要线程池

如果每来一个任务就 new Thread(),系统会同时承担线程创建销毁、线程栈内存、上下文切换和不可控并发带来的成本。线程池将“任务”和“执行资源”解耦:调用方提交 Runnable,线程池负责复用工作线程、限制最大并发、缓存突发任务并处理过载。

它主要解决四件事:

  1. 复用线程:减少频繁创建和销毁的成本。
  2. 限制并发:防止瞬时流量把 CPU、内存、数据库连接或外部服务打满。
  3. 削平峰值:用队列吸收短时间突发,但队列只能缓冲,不能创造处理能力。
  4. 统一治理:线程命名、异常处理、拒绝策略、指标监控和优雅停机可以集中管理。

需要牢记边界:JVM 里的线程池不是持久化队列。服务被强制终止后,内存队列里的任务会消失;一个 Future 超时,也不代表底层业务一定停止。

二、Executor 体系要分清

Executor
└── ExecutorService
    ├── AbstractExecutorService
    │   └── ThreadPoolExecutor
    └── ScheduledExecutorService
        └── ScheduledThreadPoolExecutor

Executors:创建 ExecutorService 的工具类,不是线程池本身
  • Executor 只有 execute(Runnable),表达“提交一个任务”。
  • ExecutorService 增加 submitshutdowninvokeAll 等生命周期和结果管理能力。
  • ThreadPoolExecutor 是最重要的通用实现,允许显式控制容量和过载行为。
  • Executors 提供快捷工厂,但快捷默认值未必符合生产环境的资源边界。

execute 适合不关心返回值的 Runnablesubmit 返回 Future,可以拿结果或异常。一个常见坑是:submit 会把任务异常封装进 Future,如果从不调用 get(),异常可能悄悄被忽略。

三、ThreadPoolExecutor 七大参数

new ThreadPoolExecutor(
    corePoolSize,
    maximumPoolSize,
    keepAliveTime,
    TimeUnit.SECONDS,
    workQueue,
    threadFactory,
    rejectedExecutionHandler
);
参数 含义 生产关注点
corePoolSize 常驻工作线程目标数量 默认也要等任务到来才创建,可选择预热线程
maximumPoolSize 线程总数上限 只有有界队列放不下任务时才可能扩到这里
keepAliveTime 超出核心数的空闲线程存活时间 配合 allowCoreThreadTimeOut 也可回收核心线程
unit 存活时间单位 与上一参数共同生效
workQueue 核心线程都忙时暂存任务 容量决定排队上限、内存风险和扩容时机
threadFactory 创建工作线程 应命名线程,并设置未捕获异常处理器
handler 饱和或关闭后继续提交时的策略 不能静默丢失关键业务任务

最容易误解的是 maximumPoolSize:如果使用无界队列,任务在核心线程忙碌后几乎总能入队,因此线程池通常不会扩到最大线程数,maximumPoolSize 基本失效。

四、execute(task) 到底怎么走

提交任务
   ↓
workerCount < corePoolSize ?
   ├─ 是:尝试创建核心工作线程执行当前任务
   └─ 否
        ↓
线程池仍在运行,并且 workQueue.offer(task) 成功?
   ├─ 是:任务入队;再次检查线程池状态
   │      ├─ 已关闭:尝试移除任务并拒绝
   │      └─ 没有工作线程:补建一个空任务 worker 来消费队列
   └─ 否
        ↓
尝试创建非核心工作线程(不超过 maximumPoolSize)
   ├─ 成功:直接执行当前任务
   └─ 失败:执行拒绝策略

这比“核心线程 → 队列 → 最大线程 → 拒绝”多了一次关键的入队后状态复检。原因是 offer 与线程池关闭可能并发发生:如果不复检,线程池关闭后仍可能遗留一个永远没人执行的任务。

源码里还有一个重要设计:线程池运行状态与 worker 数量被压进原子变量 ctl,通过 CAS 修改。这让状态判断和线程数量控制能在高并发下保持一致。

五、三种队列怎么选

队列 容量与结构 优点 风险与适用场景
ArrayBlockingQueue 固定容量、数组、FIFO 内存连续、容量清晰、性能稳定 容量创建后不可变;适合明确限制积压量的通用任务池
LinkedBlockingQueue 链表、可有界也可近似无界、FIFO 通常吞吐较高,正容量时也是 Spring 常用实现 不传容量会非常大;每个节点有对象开销,必须显式设上限
SynchronousQueue 容量为 0,只做线程间直接交接 几乎不排队,能快速扩线程处理短任务 没有消费者立刻接手就扩线程或拒绝;必须严格限制最大线程数

队列不是越大越安全。一个大队列可能降低拒绝次数,却把问题变成更长的等待时间、更旧的数据、更高的内存占用。对于 45 秒级 AI 调用,排队 100 个任务可能意味着尾部任务要等待数分钟。

六、四种拒绝策略

  1. AbortPolicy:抛出 RejectedExecutionException。失败最明显,适合关键业务,但调用方必须捕获并赋予业务语义。
  2. CallerRunsPolicy:由提交任务的线程执行,形成反压。它不丢任务,但慢任务会占住 Web 请求线程;45 秒 AI 调用通常不适合直接这样处理。
  3. DiscardPolicy:静默丢弃新任务。只适合允许丢失、且有独立指标的低价值任务。
  4. DiscardOldestPolicy:丢掉队列里最旧的任务,再尝试提交。旧任务往往更重要,业务任务通常不该使用。

真正的生产答案不是背名字,而是回答:拒绝之后,调用方看到什么?任务状态如何记录?是否会重试?重试是否幂等?是否有告警?

七、为什么不建议直接 newFixedThreadPool

Executors.newFixedThreadPool(n) 底层使用固定的核心/最大线程数和一个没有显式业务上限的 LinkedBlockingQueue。当任务到达速度长期大于消费速度时:

不会扩更多线程
    ↓
任务持续进入队列
    ↓
排队延迟持续增加
    ↓
Runnable、参数、上下文对象长期存活
    ↓
堆内存和 GC 压力增加,最终可能 OOM

因此问题不只是“一下子 OOM”,而是系统会先经历排队时间失控、超时任务仍在执行、上游重试放大流量、GC 变频繁,最后才可能内存耗尽。

同理,newCachedThreadPool 使用 SynchronousQueue 和近乎无限的最大线程数,流量失控时可能创建过多线程。快捷工厂不是语法错误,但生产系统需要明确的容量和降级策略,通常应显式构造或完整配置线程池。

八、动手实验:core=2 / max=4 / queue=2

下面的启动闸门让任务在提交阶段保持占用,因此现象更稳定:任务 1、2 创建核心线程;3、4 进入队列;5、6 创建非核心线程;7~10 被拒绝。

import java.time.LocalTime;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

public class ThreadPoolLab {
    public static void main(String[] args) throws InterruptedException {
        AtomicInteger seq = new AtomicInteger();
        CountDownLatch startGate = new CountDownLatch(1);

        ThreadFactory factory = task -> {
            Thread thread = new Thread(task);
            thread.setName("lab-worker-" + seq.incrementAndGet());
            thread.setUncaughtExceptionHandler((t, e) ->
                    System.err.println(t.getName() + " error: " + e.getMessage()));
            return thread;
        };

        ThreadPoolExecutor pool = new ThreadPoolExecutor(
                2, 4,
                5, TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(2),
                factory,
                new ThreadPoolExecutor.AbortPolicy()
        );

        for (int i = 1; i <= 10; i++) {
            int taskId = i;
            try {
                pool.execute(() -> {
                    try {
                        startGate.await();
                        System.out.printf("%s START task=%d at=%s%n",
                                Thread.currentThread().getName(), taskId, LocalTime.now());
                        TimeUnit.SECONDS.sleep(3);
                        System.out.printf("%s END   task=%d at=%s%n",
                                Thread.currentThread().getName(), taskId, LocalTime.now());
                    } catch (InterruptedException e) {
                        Thread.currentThread().interrupt();
                    }
                });
                System.out.printf("accepted task=%d pool=%d active=%d queue=%d%n",
                        taskId, pool.getPoolSize(), pool.getActiveCount(), pool.getQueue().size());
            } catch (RejectedExecutionException e) {
                System.out.printf("REJECTED task=%d%n", taskId);
            }
        }

        startGate.countDown();
        pool.shutdown();
        if (!pool.awaitTermination(30, TimeUnit.SECONDS)) {
            pool.shutdownNow();
        }
    }
}

一个容易忽略的观察:5、6 是队列满后直接交给新建的非核心线程,可能早于队列中的 3、4 完成。因此工作队列本身 FIFO,不代表整个线程池的完成顺序严格 FIFO。

将队列替换为 new LinkedBlockingQueue<>() 后再提交:大概率只有两个核心线程工作,3~10 全部入队,max=4 没有机会生效,也不会在实验规模下触发拒绝。这正是无界队列风险的直观来源。

九、线程数怎样确定

“CPU 核数 + 1”和“2 × CPU 核数”只能作为起点,不能当答案。

  • CPU 密集型:线程数通常接近可用核心数,过多线程只会增加切换。
  • I/O 密集型:可以从 N ≈ CPU核数 × 目标利用率 × (1 + 等待时间/计算时间) 估算。
  • 外部服务型任务:还必须受供应商并发/QPS、连接池、内存和成本预算约束。

对 BizeNova,更实用的并发预算是:

实例允许的 AI 并发
= min(
  AI 服务并发配额,
  QPS 配额 × 平均响应秒数,
  数据库连接预算,
  容器内存可承受并发,
  业务成本预算
)

如果部署 3 个实例,每个实例最大 8 个外层任务,总并发可能达到 24;不能只按单实例调参。还要考虑一条评估内部会并行提交多个 Agent,产生二级放大。

十、BizeNova 当前 AI 异步链路

我基于 2026-09-02 本地 main 版本梳理出的实际链路是:

PUT /companies/{companyId}/evaluations
        ↓
CompanyEvaluationServiceImpl
  权限校验 → 服务端算分 → 新建 pending 任务 → 保存访问关系
        ↓
AssessmentTaskDispatcher
  默认 local / 可选 RabbitMQ
        ↓
CompanyEvaluationAsyncProcessor
  原子 claim(pending → processing)
        ↓
私有材料上下文 + RAG 证据检索
        ↓
外部 GLM 报告生成(整体超时 45 秒)
        ↓
多个 AssessmentAgent 并行分析(单个 8 秒降级)
        ↓
company_evaluations + companies 统计 + ai_analysis_reports
        ↓
ai_evaluation_tasks → completed / failed
        ↓
SSE 进度与轮询接口返回结果

1. 两组线程池

// 外层评估任务池
corePoolSize = 4;
maxPoolSize = 8;
queueCapacity = 100;

// 一条评估内部的 Agent 池
corePoolSize = 5;
maxPoolSize = 10;
queueCapacity = 200;

Spring 在正数 queueCapacity 下创建有界 LinkedBlockingQueue。两组池分离是优点:外层任务会等待 Agent 结果,如果共用一个小池,所有线程都可能在等待自己提交的子任务,形成线程饥饿。问题是当前参数硬编码,缺少压测依据、拒绝计数和排队耗时监控。

2. 原子领取避免同 taskId 重复消费

UPDATE ai_evaluation_tasks
SET task_status = 'processing', start_time = NOW(), update_time = NOW()
WHERE id = #{taskId}
  AND task_status = 'pending'
  AND deleted = 0;

只有一个 worker 能把受影响行数更新为 1,其余重复消息会拿到 0 并退出。这是正确的幂等入口,但只解决“同一个 taskId 同时执行”,不能解决“客户端重试创建了两个不同 taskId”。

3. AI 超时与降级

当前 HTTP 客户端设置 10 秒连接超时,单次报告请求整体上限 45 秒。非 2xx、超时、连接异常、响应过大、JSON 不合法或输出被截断时,不再次付费重试,而是返回明确标记 _degraded=true 的保守报告。

这是一个合理的安全兜底:评估不会因为模型不可用无限挂住,也不会伪造精确结论。但任务最终仍被写成 completed,监控层必须把“正常完成”和“降级完成”分开,否则供应商故障会被成功率掩盖。

4. Agent 并行不等于真正取消

CompletableFuture.supplyAsync(callAgent, assessmentAgentExecutor)
    .orTimeout(8, TimeUnit.SECONDS)
    .exceptionally(error -> conservativeFallback);

orTimeout 能让等待方在 8 秒后得到降级结果,但通常不能保证正在运行的 supplier 立刻停止。如果内部网络调用没有自己的超时或不响应中断,线程仍可能继续被占用。因此“Future 设置超时”和“底层 I/O 可取消”必须分别设计。

十一、七个项目问题的答案

1. 线程池参数是多少?

外层 4/8/100,Agent 层 5/10/200。两者都使用有界队列。

2. 为什么这样配置?

当前代码没有给出压测报告或容量模型,所以不能声称这些值已经合理。需要采集 AI P50/P95 延迟、单任务内存、连接池占用、到达率与排队 SLO 后决定。

3. 队列满了怎么办?

没有显式设置 handler,默认拒绝并抛异常。此时任务已经写入数据库,但接口可能返回 500;客户端如果重试,就可能创建第二个任务。慢 AI 不建议直接改成 CallerRunsPolicy,否则 Tomcat 请求线程会同步执行几十秒。

4. AI 请求超时怎么办?

当前整体 45 秒后降级,并生成保守报告;并行 Agent 单个 8 秒后降级。还要确保每个底层 HTTP/数据库调用都有自己的超时。

5. AI API 返回 500 怎么办?

当前直接降级,不重试。优点是不会造成重试风暴和重复计费;缺点是瞬时 500 也无法恢复。改进时应只对 429、部分 5xx、连接重置做有限次数指数退避,并遵守 Retry-After

6. 服务执行一半重启会不会丢?

会有风险。系统每分钟扫描创建超过 2 分钟的 pending 任务并重投,但已经 claim 为 processing 的任务不在扫描范围内。进程中途退出后,它可能永久卡在 processing。

7. 同一个任务重复执行会不会产生脏数据?

原子 claim 能挡住同 taskId 的并发重复执行;但评估记录、公司统计、报告和任务状态是多次独立写入。如果写完评估后、写报告前失败,前半段已经提交。以后人工重置任务再跑,会重复插入评估并重复累加统计。结果表缺少 task_id 唯一约束,这是当前最关键的数据一致性风险之一。

十二、至少六个可靠性风险

  1. processing 没有租约、心跳和超时恢复,重启后可能永久卡住。
  2. 多表写入缺少明确事务或阶段级幂等,部分成功后重试可能产生脏数据。
  3. 创建接口缺少 request_key,网络重试可能创建多个业务等价任务。
  4. 线程池参数硬编码,没有 active、queue、reject、wait time 等指标。
  5. 默认拒绝异常缺少“任务已落库、稍后恢复”的业务响应设计。
  6. orTimeout 不保证底层 Agent 停止,慢调用可能继续占用线程。
  7. 降级报告与完整报告都归入 completed,容易误判业务成功率。
  8. 任务级线程池与 Agent 线程池形成并发放大,必须统一做实例级预算。

十三、暂不上 MQ,先把数据库任务做可靠

当前已经有 ai_evaluation_tasks,完全可以先把它升级为可靠任务事实源:

新增字段
request_key       唯一请求幂等键
retry_count       已重试次数
max_retries       最大重试次数
next_retry_at     下次可执行时间
worker_id         当前执行者
lease_until       执行租约截止时间
heartbeat_at      最近心跳
result_quality    normal / degraded
version           乐观锁版本

推荐状态机:

pending → processing → completed
              ├──────→ completed_degraded
              ├──────→ retry_wait → pending
              └──────→ failed

processing 且 lease_until 过期 → retry_wait

实现顺序:

  1. 创建任务与访问记录使用短事务,request_key 唯一。
  2. 线程池拒绝时保留 pending,接口返回已受理 taskId,由调度器稍后领取。
  3. 多实例通过条件更新或 FOR UPDATE SKIP LOCKED 领取,并设置租约。
  4. 长任务刷新 heartbeat,watchdog 回收过期租约。
  5. 按错误类型决定重试,使用指数退避和随机抖动。
  6. company_evaluations.task_idai_analysis_reports.task_id 建唯一索引,重复执行使用 upsert。
  7. 暴露任务成功率、降级率、重试率、P95 排队时间、P95 执行时间和线程池饱和度。

十四、线程池优化代码方向

下面是“配置化 + 显式拒绝 + 优雅停机”的方向示例,数值仍需压测:

@Bean("assessmentTaskExecutor")
public ThreadPoolTaskExecutor assessmentTaskExecutor(AssessmentPoolProperties p) {
    ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
    executor.setCorePoolSize(p.coreSize());
    executor.setMaxPoolSize(p.maxSize());
    executor.setQueueCapacity(p.queueCapacity());
    executor.setKeepAliveSeconds(p.keepAliveSeconds());
    executor.setThreadNamePrefix("assessment-");
    executor.setRejectedExecutionHandler(new ThreadPoolExecutor.AbortPolicy());
    executor.setWaitForTasksToCompleteOnShutdown(true);
    executor.setAwaitTerminationSeconds(30);
    executor.initialize();
    return executor;
}

为什么仍选择显式 AbortPolicy?因为关键任务不该静默丢弃,而 45 秒 AI 调用也不适合占用请求线程。正确做法是在 dispatch 边界捕获拒绝,记录 reject 指标,保留 pending 任务,让恢复调度器稍后处理。

十五、什么时候再引入 RabbitMQ / RocketMQ

BizeNova 代码里已经存在可选 RabbitMQ dispatcher、持久化队列和 DLQ,但默认开关为 false。今天不需要为了展示技术栈就启用它。

当出现多实例可靠分发、明显削峰需求、消费者独立扩缩容、延迟重试/死信治理,或数据库轮询已经成为瓶颈时,再让 MQ 成为主链路。引入时最重要的不是换一个 send()

数据库事务:task + outbox
        ↓
Outbox Publisher + Publisher Confirm
        ↓
MQ
        ↓
Consumer 原子 claim
        ↓
幂等持久化成功后 ACK
        ↓
有限重试 → DLQ → 告警/人工补偿

Transactional Outbox 用来解决“数据库任务创建成功,但消息没发出去”或“消息发出去了,但数据库事务回滚”的双写一致性。无论以后使用 RabbitMQ 还是 RocketMQ,数据库任务状态机、租约和幂等约束都仍然有价值。

十六、企业中的实际应用场景

适合进程内线程池的场景:

  • 一次请求内并行调用多个相互独立的查询,最后汇总结果;
  • 可重建、可丢失或已有数据库事实源兜底的后台任务;
  • 图片处理、报表局部计算、批量数据分片;
  • 邮件/通知投递前的轻量组装;
  • 调用外部 AI、OCR、搜索服务时限制单实例并发。

不应只依赖内存线程池的场景:

  • 订单扣款、库存扣减等不能丢失的核心交易;
  • 跨服务长流程,必须重试、追踪和补偿;
  • 任务执行时间很长,可能跨越发布重启;
  • 多实例需要全局顺序、延迟消息或统一削峰;
  • 业务要求精确审计每次尝试与结果。

线程池、数据库任务表和 MQ 不是互相替代,而是不同层次:线程池控制单进程执行资源;数据库状态机保证业务可追踪与幂等;MQ 负责跨实例可靠传输、削峰和解耦。

十七、算法:LeetCode 146 LRU 缓存

LRU 要求 getput 都是 O(1)。单独用 HashMap 能 O(1) 查找,却不能 O(1) 找到最久未使用元素;单独用链表能维护顺序,却不能 O(1) 定位节点。因此组合为:

HashMap<key, Node>:O(1) 定位
双向链表:O(1) 删除节点、移动到头部、淘汰尾部
import java.util.HashMap;
import java.util.Map;

class LRUCache {
    private static class Node {
        int key, value;
        Node prev, next;
        Node(int key, int value) { this.key = key; this.value = value; }
    }

    private final int capacity;
    private final Map<Integer, Node> cache = new HashMap<>();
    private final Node head = new Node(0, 0); // 最近使用端
    private final Node tail = new Node(0, 0); // 最久未使用端

    LRUCache(int capacity) {
        this.capacity = capacity;
        head.next = tail;
        tail.prev = head;
    }

    public int get(int key) {
        Node node = cache.get(key);
        if (node == null) return -1;
        moveToHead(node);
        return node.value;
    }

    public void put(int key, int value) {
        Node node = cache.get(key);
        if (node != null) {
            node.value = value;
            moveToHead(node);
            return;
        }
        Node fresh = new Node(key, value);
        cache.put(key, fresh);
        addAfterHead(fresh);
        if (cache.size() > capacity) {
            Node victim = tail.prev;
            remove(victim);
            cache.remove(victim.key);
        }
    }

    private void moveToHead(Node node) {
        remove(node);
        addAfterHead(node);
    }

    private void addAfterHead(Node node) {
        node.prev = head;
        node.next = head.next;
        head.next.prev = node;
        head.next = node;
    }

    private void remove(Node node) {
        node.prev.next = node.next;
        node.next.prev = node.prev;
    }
}

哨兵节点 head/tail 消除了插入空链表、删除首尾节点的分支。面试时要主动说明:这份实现不是线程安全的;如果并发访问,需要外部锁、分段设计或直接使用成熟缓存组件。

十八、高频面试题与参考答案

Q1:ThreadPoolExecutor 七个参数是什么?

核心线程数、最大线程数、非核心线程空闲存活时间、时间单位、工作队列、线程工厂、拒绝策略。回答时要补充参数之间的联动,尤其是队列是否有界会决定最大线程数能否生效。

Q2:任务提交后的执行流程?

先尝试补到核心线程数;核心线程满后尝试入队;入队后复检运行状态;队列满后再扩到最大线程数;仍无法接收时执行拒绝策略。

Q3:核心线程一定在创建线程池时就启动吗?

不一定。默认通常在任务到来时按需创建,可用 prestartCoreThreadprestartAllCoreThreads 预热。

Q4:maximumPoolSize 什么时候会失效?

使用几乎无界的工作队列时,核心线程满后的任务一直能入队,线程池不会走“队列满后扩线程”的分支,因此最大线程数基本不起作用。

Q5:为什么生产环境强调有界队列?

有界队列把过载变成可观测、可处理的拒绝或降级;无界队列把压力藏进内存和排队延迟,最终可能造成超时放大、GC 压力与 OOM。

Q6:四种拒绝策略如何选?

关键任务通常以 AbortPolicy 为基础,捕获异常并做持久化重试或快速失败;CallerRunsPolicy 能反压但会占调用线程;两个 Discard 策略只适合明确允许丢失的任务,并必须有指标。

Q7:execute 与 submit 有什么区别?

execute 接收 Runnable、无返回值,任务抛出的未捕获异常可到线程异常处理器;submit 返回 Future,异常被封装,需要 get() 或统一 afterExecute 处理才能观察。

Q8:线程数怎么估算?

CPU 密集接近核心数;I/O 密集可按等待/计算比估算,但最终要受下游 QPS、连接池、内存和多实例总并发约束,并通过压测和 P95 指标调整。

Q9:Future 超时后任务一定停了吗?

不一定。超时可能只让等待方结束;底层任务是否停止取决于取消是否传递、线程是否被中断、业务代码是否响应中断、网络库是否支持取消以及是否设置自身超时。

Q10:@Async 有什么常见坑?

同类内部调用绕过代理会导致异步失效;未指定线程池可能使用不合适的默认执行器;void 方法异常不返回给调用方;线程上下文、事务和 MDC 不会天然完整传播;进程退出会丢失内存队列任务。

Q11:BizeNova 为什么需要异步 AI 评估?

外部 AI 与 RAG 延迟高且波动大,同步等待会长时间占用请求线程,增加网关超时概率。异步后接口快速返回 taskId,后台限并发执行,前端通过 SSE/轮询看进度,能隔离外部服务波动。

Q12:如果重新设计 BizeNova,你会怎么做?

先保留数据库任务表,不急着上 MQ:增加请求幂等键、任务租约、心跳、重试时间和结果质量;结果表按 taskId 唯一并 upsert;线程池参数配置化、指标化;对可重试错误做有限退避。多实例和任务量上升后,再用 Outbox + MQ 做可靠投递,消费者仍通过原子 claim 和幂等写入保证至少一次投递下的正确性。

Q13:BizeNova 当前至少三个风险是什么?

processing 任务重启后不能恢复;多表分步写导致部分成功和重复写;缺少请求级幂等;另外还有默认拒绝缺少业务处理、线程池不可观测、Future 超时不保证取消等风险。

Q14:为什么不能只把队列调大?

队列变大不会提高消费能力,只会推迟拒绝并增加等待时间和内存占用。对于有时效性的 AI 评估,等待过久的结果可能已经失去价值。容量应由允许排队时长和任务到达/处理速率反推。

十九、今天的闭卷检查清单

完成今天学习后,应能不看资料说清:

  • 七大参数及其联动关系;
  • execute 的四段决策与入队后复检;
  • 三种典型队列的容量模型;
  • 四种拒绝策略的业务后果;
  • 无界队列为什么让 maximumPoolSize 失效;
  • 线程数为什么必须结合下游容量与压测;
  • BizeNova 两组线程池分别处理什么;
  • 当前 pending 能恢复、processing 不能恢复的原因;
  • 原子 claim 能解决什么、不能解决什么;
  • 不上 MQ 时如何用任务表、租约和幂等做到可靠异步;
  • LRU 为什么使用 HashMap + 双向链表。

总结

线程池不是“让代码异步”的注解,而是一套容量管理机制。真正的项目亮点也不是写出 @Async,而是能解释:为什么分池、并发上限来自哪里、队列满了如何反馈、超时任务是否真正停止、服务重启如何恢复、重复执行如何幂等、什么时候才值得引入 MQ。

BizeNova 已经具备任务表、原子 claim、pending 恢复、AI 超时降级、SSE 进度以及可选 RabbitMQ 等不错的基础。下一步最有价值的不是更换中间件,而是补齐 processing 租约恢复、请求幂等、阶段幂等、事务边界和可观测性。把这些问题讲清楚,线程池知识才真正从八股变成项目经验。

参考资料