在线判题系统的并发提交处理:消息队列解耦与异步回调设计
在线判题系统的并发提交处理消息队列解耦与异步回调设计一、深度引言与场景痛点当 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 ApiResponseSubmitResponse 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() { // 从死信队列中拉取消息记录到告警表 // 这些是自动重试也无法恢复的异常需要人工排查 ListJudgeTaskMessage 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 层和判题层幂等性是异步系统的必修课——永远假设消息可能被投递多次死信队列是最后的安全网——不要让它悄无声息地丢消息对于实习生来说把一个同步服务拆成异步架构是理解分布式系统设计最好的入门项目。因为你能亲手感受到解耦之后的系统调试难度也在同步上升。