将催款通知、AI 问答、普通操作分离到不同的独立队列,各自使用独立的自定义线程池处理。
@Component
public class PriorityTaskDispatcher {
// 高优先级队列 —— AI 问答(独立线程池)
private final BlockingQueue<Runnable> highPriorityQueue = new LinkedBlockingQueue<>(100);
private final ExecutorService highExecutor = new ThreadPoolExecutor(
3, 5, 60L, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(50),
new ThreadFactoryBuilder().setNameFormat("ai-pool-%d").build(),
new ThreadPoolExecutor.CallerRunsPolicy()
);
// 低优先级队列 —— 群发催款(带限流)
private final BlockingQueue<Runnable> lowPriorityQueue = new LinkedBlockingQueue<>(1000);
private final ExecutorService lowExecutor = Executors.newFixedThreadPool(5);
private volatile int batchDelayMs = 100; // 每个催款消息间隔 100ms
public void setBatchDelay(int ms) { this.batchDelayMs = ms; }
// 批量任务执行时加入延迟,控制发送频率
public void postBatchTask(Runnable task) {
lowPriorityQueue.offer(() -> {
try {
task.run();
Thread.sleep(batchDelayMs); // 控制发送频率,避免手机被封控
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
});
}
}
TaskDispatcher,业务层代码基本不变。batchDelayMs 控制发送频率,防止账号被风控。基于方案一的缺陷,选择 RabbitMQ 作为最终方案,利用其天然的顺序保证、持久化、优先级等特性。
@Configuration
public class RabbitMQConfig {
@Bean
public Queue aiQueue() {
return QueueBuilder.durable("ai.queue")
.maxPriority(10) // 支持优先级 0-10
.deadLetterExchange("dlx.exchange")
.deadLetterRoutingKey("dlx.routing")
.build();
}
@Bean
public Queue batchQueue() {
return QueueBuilder.durable("batch.queue")
.maxPriority(5)
.build();
}
@Bean
public TopicExchange exchange() {
return new TopicExchange("task.exchange");
}
}
@Component
@Slf4j
public class TaskProducer {
@Autowired
private RabbitTemplate rabbitTemplate;
public void sendTask(MessageDTO dto) {
String routingKey = isAiTask(dto) ? "ai.task" : "batch.task";
int priority = isAiTask(dto) ? 10 : 1; // AI 问答最高优先级
rabbitTemplate.convertAndSend("task.exchange", routingKey, dto,
msg -> {
msg.getMessageProperties().setPriority(priority);
return msg;
});
}
private boolean isAiTask(MessageDTO dto) {
return "发文本".equals(dto.getAct())
&& dto.getMsg() != null
&& dto.getMsg().contains("?");
}
}
@Component
@Slf4j
public class TaskConsumer {
@RabbitListener(queues = "ai.queue", concurrency = "3-5") // 动态伸缩
public void handleAiTask(MessageDTO dto, Channel channel, Message message) {
try {
actWithPhone(dto);
channel.basicAck(message.getMessageProperties().getDeliveryTag(), false);
} catch (Exception e) {
log.error("AI 任务执行失败,消息ID: {}",
message.getMessageProperties().getMessageId(), e);
// 重试(requeue=true),失败后进入死信队列
channel.basicNack(message.getMessageProperties().getDeliveryTag(), false, true);
}
}
}
因此,当前系统采用 RabbitMQ 作为任务调度核心组件,既解决了单线程瓶颈,又为未来业务增长预留了扩展空间。