在线判题系统的并发提交处理:消息队列解耦与异步回调设计
一、深度引言与场景痛点:当 100 个人同时点"提交",系统发生了什么?
V1 版本的判题系统用的是最简单的同步模式:用户点"提交" → HTTP 请求进入 → 编译代码 → 运行测试用例 → 返回结果 → 用户看到结果。这个流程在单人使用时完美运行,响应速度也很快。
直到有一次内部测试,10 个同事同时提交了不同的题目。崩溃发生了:Tomcat 的请求线程池被打满,新的请求在排队,旧的请求因为判题时间过长把线程一直占着。最终,整个服务不可用,连健康检查接口都 503 了。
问题的根源很简单:判题是一个长耗时、CPU 密集型的操作,但它被绑定在了短生命周期的 HTTP 请求线程上。HTTP 请求线程的职责是接收请求、返回响应,它不应该被一个可能要跑 5 秒的判题过程阻塞。
二、底层机制与原理深度剖析
异步解耦的核心思想:把"提交判题"和"执行判题"拆开。
这套架构的关键收益:
- API 服务不再阻塞:提交接口的响应时间从秒级降到毫秒级
- 判题 Worker 独立扩缩:可以通过增加 Worker 数量来提升并发处理能力
- 故障隔离:判题 Worker 崩溃不会影响 API 服务的可用性
- 削峰填谷:消息队列天然具备缓冲能力,应对瞬时提交高峰
三、生产级代码实现与最佳实践
消息结构定义
// 判题任务消息体 —— 包含判题所需的所有信息,Worker 可独立处理 @Data @AllArgsConstructor @NoArgsConstructor public class JudgeTaskMessage implements Serializable { private String submissionId; // 提交记录 ID private String problemId; // 题目 ID private String code; // 用户代码 private String language; // 编程语言 private Long submitTime; // 提交时间戳 // 最大重试次数 —— 防止逻辑错误导致无限重试 private int maxRetryCount = 3; private int currentRetryCount = 0; }提交接口实现
// 提交接口 —— 只负责接收和记录,不执行判题 @RestController @RequestMapping("/api/submission") public class SubmissionController { private final SubmissionService submissionService; private final RabbitTemplate rabbitTemplate; @PostMapping("/submit") public ApiResponse<SubmitResponse> submit(@RequestBody SubmitRequest request) { // 1. 创建提交记录,状态设为 PENDING // 即使判题失败,记录也已经持久化,用户可以追溯 Submission submission = submissionService.createPendingSubmission( request.getProblemId(), request.getCode(), request.getLanguage() ); // 2. 构造消息并发送到队列 // 使用 convertAndSend 确保消息序列化后进入队列 JudgeTaskMessage message = new JudgeTaskMessage( submission.getId(), request.getProblemId(), request.getCode(), request.getLanguage(), System.currentTimeMillis() ); // exchange: judge.exchange // routingKey: judge.task.{language} // 按语言路由可以让不同语言的 Worker 分别消费 String routingKey = "judge.task." + request.getLanguage().toLowerCase(); rabbitTemplate.convertAndSend("judge.exchange", routingKey, message); // 3. 立即返回 submissionId,用户用它轮询结果 return ApiResponse.success(new SubmitResponse(submission.getId())); } }判题 Worker 实现
// 判题 Worker —— 独立消费消息,专注执行判题逻辑 @Component @Slf4j public class JudgeWorker { private final SubmissionRepository submissionRepository; private final JudgeService judgeService; // 并发消费者数量 —— 通过配置中心动态调整 // concurrency: 消费者线程数,对应同时判题的并发数 @RabbitListener( queues = "judge.queue.java", concurrency = "3-10" // 最少 3 个,最多 10 个消费者 ) public void handleJudgeTask(JudgeTaskMessage message) { log.info("开始判题: submissionId={}, language={}", message.getSubmissionId(), message.getLanguage()); try { // 1. 更新状态为 RUNNING,用户端可以看到"判题中" submissionRepository.updateStatus( message.getSubmissionId(), SubmissionStatus.RUNNING ); // 2. 执行实际判题逻辑 JudgeResult result = judgeService.judge( message.getProblemId(), message.getCode(), message.getLanguage() ); // 3. 更新最终结果 submissionRepository.updateResult( message.getSubmissionId(), result.getStatus(), result.getDetails() ); log.info("判题完成: submissionId={}, result={}", message.getSubmissionId(), result.getStatus()); } catch (Exception e) { log.error("判题异常: submissionId={}", message.getSubmissionId(), e); handleJudgeFailure(message, e); } } private void handleJudgeFailure(JudgeTaskMessage message, Exception e) { int retryCount = message.getCurrentRetryCount() + 1; if (retryCount <= message.getMaxRetryCount()) { // 重试机制:更新重试计数,重新投递到延迟队列 // 延迟 30 秒后重试,给系统恢复的时间窗口 message.setCurrentRetryCount(retryCount); rabbitTemplate.convertAndSend( "judge.exchange", "judge.retry", message, msg -> { // 设置消息的 TTL 实现延迟投递 msg.getMessageProperties().setExpiration("30000"); return msg; } ); } else { // 超过最大重试次数,标记为系统错误 submissionRepository.updateStatus( message.getSubmissionId(), SubmissionStatus.SYSTEM_ERROR ); log.error("判题重试耗尽: submissionId={}, retryCount={}", message.getSubmissionId(), retryCount); } } }死信队列处理
// 死信队列配置 —— 处理多次重试后仍然失败的消息 @Configuration public class DeadLetterConfig { // 死信队列:接收重试耗尽的消息,避免消息丢失 @Bean public Queue deadLetterQueue() { return QueueBuilder.durable("judge.dlq").build(); } @Bean public Binding deadLetterBinding() { return BindingBuilder .bind(deadLetterQueue()) .to(deadLetterExchange()) .with("judge.dlq"); } // 定时任务:每天凌晨检查死信队列,人工介入处理 @Scheduled(cron = "0 0 2 * * ?") public void processDeadLetters() { // 从死信队列中拉取消息,记录到告警表 // 这些是自动重试也无法恢复的异常,需要人工排查 List<JudgeTaskMessage> deadMessages = fetchDeadMessages(); if (!deadMessages.isEmpty()) { log.warn("死信队列中有 {} 条未处理消息,请尽快排查", deadMessages.size()); alertService.sendAlert("判题死信队列堆积", deadMessages.size()); } } }四、边界分析与架构权衡
消息丢失问题
RabbitMQ 的消息持久化 + 手动确认(Manual ACK)可以最大程度防止消息丢失,但以下情况仍需要额外处理:
- Worker 在处理中崩溃:消息会重新入队(如果开启了 ACK),由下一个 Worker 继续处理
- RabbitMQ 本身宕机:磁盘上的持久化消息可以恢复,但内存中的瞬态消息会丢失
- 网络分区:使用镜像队列(Mirrored Queue)或 Quorum Queue 提高可用性
补偿方案:在 API 层增加定时扫描任务,检查超过 N 分钟仍为 PENDING 状态的提交记录,主动补发判题消息。
幂等性保证
同一个提交可能被多次投递(网络重试、Worker 重连等)。Worker 端必须保证判题的幂等性:
// 幂等性检查 —— 提交 ID 作为幂等键 if (submissionRepository.existsById(message.getSubmissionId())) { Submission existing = submissionRepository.findById(message.getSubmissionId()); if (existing.getStatus() != SubmissionStatus.PENDING) { log.info("跳过重复判题: submissionId={}, 当前状态={}", message.getSubmissionId(), existing.getStatus()); return; // 已经处理过了,直接返回 } }同步 vs 异步的选择时机
| 系统阶段 | 推荐方案 | 原因 |
|---|---|---|
| MVP(< 10 用户) | 同步判题 | 实现简单,快速验证 |
| 内测(< 100 用户) | 异步 + 单 Worker | 解耦接口,预留扩展空间 |
| 正式运营(100+ 用户) | 异步 + 多 Worker + MQ | 高并发下的唯一选择 |
五、总结
将判题从同步改为异步,是这个系统架构变化最大的一次升级。核心经验有几点:
- 不要把长耗时操作绑定在 HTTP 线程上——它们是稀缺资源
- 消息队列的解耦作用远大于削峰——它让你可以独立演进 API 层和判题层
- 幂等性是异步系统的必修课——永远假设消息可能被投递多次
- 死信队列是最后的安全网——不要让它悄无声息地丢消息
对于实习生来说,把一个同步服务拆成异步架构,是理解"分布式系统设计"最好的入门项目。因为你能亲手感受到:解耦之后的系统,调试难度也在同步上升。