本周的任务仍然聚焦于错题本模块AI系统的开发。本周的任务侧重于修改和完善之前的一些历史遗留问题和开发业务端和AI端的消息队列信息传递的相关代码

一.解决历史遗留问题:

考虑到本系统需要有用户端与管理员端的相关功能,之前的用户端设计有缺陷,即不能够很好地兼容用户端与管理员端,所以本周经过审议,我对管理员端的数据库表进行了修改,添加了部分属性:

我为user表添加了status和role属性,status表示用户当前的状态(是否被封禁),role表示角色(用户或者管理员),并且同时修改相关的代码:

@TableName("user")
@Data
@Schema(description = "用户实体类")
public class User {
//修改用户实体类
    @Schema(description = "主键ID")
    @TableId(type = IdType.ASSIGN_ID)
    private Long id;

    @Schema(description = "用户名")
    private String name;

    @Schema(description = "密码")
    private String password;

    @Schema(description = "头像URL")
    private String avatar;

    @Schema(description = "创建时间")
    private LocalDateTime createTime;

    @Schema(description = "更新时间")
    private LocalDateTime updateTime;

    @Schema(description = "逻辑删除标识:0-未删除,1-已删除")
    private Integer deleted;

    @Schema(description = "账号状态:0-正常,1-封禁")
    private Integer status;

    @Schema(description = "角色:0-普通用户,1-管理员")
    private Integer role;
}
@Data
@Builder
public class LoginVO {

    private Long id;
    private String name;

    //private String password;  //修改_0418: 不返回密码
    private String avatar;

    //新添加角色信息
    private Integer role;
    private String token;
}
//在登录和注册方法中,添加了判断状态的代码以及设置角色的代码 
@Override
    public User login(LoginDTO loginDTO) {


        LambdaQueryWrapper<User> queryWrapper = new LambdaQueryWrapper<>();
        queryWrapper.eq(User::getName, loginDTO.getName());
        queryWrapper.eq(User::getPassword, loginDTO.getPassword());
        queryWrapper.eq(User::getDeleted, 0);
        queryWrapper.eq(User::getStatus, 1);
        User user = getBaseMapper().selectOne(queryWrapper);
        return user;
    }

    @Override
    public void register(RegisterDTO registerDTO) {
        LambdaQueryWrapper<User> queryWrapper = new LambdaQueryWrapper<>();
        queryWrapper.eq(User::getName, registerDTO.getName());
        User validUser = getBaseMapper().selectOne(queryWrapper);
        if(validUser != null) {
            throw new RuntimeException("用户已存在");
        }
        User user = new User();
        user.setName(registerDTO.getName());
        user.setPassword(registerDTO.getPassword());
        user.setAvatar(registerDTO.getAvatar());
        user.setCreateTime(LocalDateTime.now());
        user.setUpdateTime(LocalDateTime.now());

        user.setDeleted(0);
        user.setRole(0);
        user.setStatus(1);
        save(user);

    }

二.设计和完善与RabbitMQ交互的工具类

在错题分析相关的AI模块,设计这两个队列,一个队列负责发送数据,一个队列负责接收数据。在AI端则是相反的方向。

// 业务端发送数据的队列
    public static final String PICTURE_QUEUE = "picture-routing-key";
    // AI分析结果队列(接收AI返回的结果)
    public static final String AI_RESULT_QUEUE = "ai-result-routing-key";

设计相关的交互代码,用来发送数据和消费数据,此为发送数据的工具类:

@Slf4j
@Component
public class RabbitMQUtils {

    @Resource
    private RabbitTemplate rabbitTemplate;

    /**
     * 发送消息(同步)
     * 
     * @param routingKey 路由键(对应Kafka的topic)
     * @param message 消息内容
     * @return 是否发送成功
     */
    public boolean sendSync(String routingKey, Object message) {
        try {
            // 使用CompletableFuture实现同步等待
            CompletableFuture<Boolean> future = new CompletableFuture<>();
            
            CorrelationData correlationData = new CorrelationData(UUID.randomUUID().toString());
            
            rabbitTemplate.convertAndSend(RabbitMQConfig.EXCHANGE_NAME, routingKey, message, correlationData);
            
            // 简单模拟确认(实际生产环境需要配置publisher confirms)
            future.complete(true);
            log.info("消息发送成功 - RoutingKey: {}, Message: {}", routingKey, message);
            return true;
        } catch (Exception e) {
            log.error("消息发送失败 - RoutingKey: {}, Message: {}", routingKey, message, e);
            return false;
        }
    }

    /**
     * 发送消息(异步,带回调)
     * 
     * @param routingKey 路由键(对应Kafka的topic)
     * @param message 消息内容
     */
    public void sendAsync(String routingKey, Object message) {
        CompletableFuture.runAsync(() -> {
            try {
                rabbitTemplate.convertAndSend(RabbitMQConfig.EXCHANGE_NAME, routingKey, message);
                log.info("消息异步发送成功 - RoutingKey: {}, Message: {}", routingKey, message);
            } catch (Exception e) {
                log.error("消息异步发送失败 - RoutingKey: {}, Message: {}", routingKey, message, e);
            }
        });
    }

此为消费数据的工具类,目前为消费ai-result-routing-key

//类  RabbitMQConsumerService
@RabbitListener(queues = "ai-result-routing-key",containerFactory = "rabbitListenerContainerFactory")
    public void consumeAIResult(String message, Channel channel,
                                @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) {
        log.info("接收到AI分析结果 - Queue: ai-result-queue, Message: {}", message);
        
        try {
            // 1. 校验消息
            if (message == null || message.trim().isEmpty()) {
                log.warn("收到空消息,丢弃并确认");
                channel.basicAck(deliveryTag, false);
                return;
            }
            
            // 2. 校验JSON格式
            String trimmedMessage = message.trim();
            if (!trimmedMessage.startsWith("{")) {
                log.error("AI结果格式错误:期望 JSON 对象,实际内容: {}", message);
                channel.basicNack(deliveryTag, false, false);
                return;
            }

            log.info("AI结果已接收,开始解析 - Content: {}", trimmedMessage);
            // 3. 解析AI分析结果
            TopicAnalysePackage analysePackage = JSON.parseObject(trimmedMessage, TopicAnalysePackage.class);
            
            if (analysePackage == null) {
                log.error("AI结果解析为 null - Content: {}", message);
                channel.basicNack(deliveryTag, false, false);
                return;
            }
            
            log.info("AI分析结果解析成功 - RequestId: {}, QuestionStem: {}", 
                    analysePackage.getRequestId(), analysePackage.getQuestionStem());
            
            // 4. 通知等待的请求(关键步骤!)
            if (analysePackage.getRequestId() != null) {
                aiResultManager.completeRequest(analysePackage.getRequestId(), analysePackage);
                log.info("已通知等待的请求 - RequestId: {}", analysePackage.getRequestId());
            } else {
                log.warn("AI结果中缺少 requestId,无法通知等待方");
            }
            

            
            log.info("AI分析结果处理完成");
            
            // 6. 手动确认消息
            channel.basicAck(deliveryTag, false);
            
        } catch (Exception e) {
            log.error("AI分析结果处理失败 - Message: {}", message, e);
            try {
                channel.basicNack(deliveryTag, false, false);
            } catch (IOException ioException) {
                log.error("拒绝消息时发生异常", ioException);
            }
        }
    }

设置AI消息管理类,用于暂存队列中的消息,并设置时限(180s)

public class AIResultManager {

    // 存储等待中的请求: requestId -> CompletableFuture
    private final Map<String, CompletableFuture<TopicAnalysePackage>> pendingRequests = new ConcurrentHashMap<>();

    /**
     * 创建一个新的等待请求
     * 
     * @param requestId 请求ID
     * @return CompletableFuture,调用方可以等待它获取结果
     */
    public CompletableFuture<TopicAnalysePackage> createPendingRequest(String requestId) {
        CompletableFuture<TopicAnalysePackage> future = new CompletableFuture<>();
        pendingRequests.put(requestId, future);
        
        log.info("创建AI分析等待请求 - RequestId: {}", requestId);
        
        // 设置超时自动清理(180秒)
        future.orTimeout(180, TimeUnit.SECONDS)
              .exceptionally(ex -> {
                  log.warn("AI分析请求超时或被取消 - RequestId: {}", requestId);
                  pendingRequests.remove(requestId);
                  return null;
              });
        
        return future;
    }

实现接口uploadAI的服务类方法:

 @Override
    public TopicAnalyseVO uploadAI(BaseRequest baseRequest) {
        log.info("开始处理AI图片分析请求: {}", baseRequest);
        
        try {
            // 1. 从 JWT Token 中获取当前用户ID
            Long userId = BaseContext.getCurrentId();
            if (userId == null) {
                log.error("未获取到用户ID,无法进行AI分析");
                throw new RuntimeException("用户未登录");
            }
            
            // 2. 生成唯一的请求ID
            String requestId = UUID.randomUUID().toString();
            
            // 3. 创建等待请求(返回 CompletableFuture)
            CompletableFuture<TopicAnalysePackage> resultFuture = aiResultManager.createPendingRequest(requestId);
            
            // 4. 构建 UrlPackage 对象
            UrlPackage urlPackage = new UrlPackage();
            urlPackage.setUrl(baseRequest.getId()); // id 字段存储的是图片URL
            urlPackage.setUserId(userId);
            urlPackage.setRequestId(requestId); // 设置请求ID
            
            // 5. 将 UrlPackage 转换为 JSON 字符串
            String jsonMessage = JSON.toJSONString(urlPackage);
            
            // 6. 发送到 picture-queue 队列
            rabbitMQUtils.sendAsync(RabbitMQConfig.PICTURE_QUEUE, jsonMessage);
            
            log.info("图片URL已发送到RabbitMQ - UserId: {}, RequestId: {}, Url: {}", 
                    userId, requestId, baseRequest.getId());
            
            // 7. 同步等待AI处理结果(最多等待30秒)
            TopicAnalysePackage analysePackage = resultFuture.get(180, TimeUnit.SECONDS);
            
            if (analysePackage == null) {
                log.error("AI分析结果为空 - RequestId: {}", requestId);
                throw new RuntimeException("AI分析失败,未返回结果");
            }
            
            log.info("AI分析完成 - RequestId: {}, QuestionStem: {}", 
                    requestId, analysePackage.getQuestionStem());
            
            // 8. 转换为 VO 返回
            TopicAnalyseVO resultVO = new TopicAnalyseVO();
            resultVO.setQuestionStem(analysePackage.getQuestionStem());
            resultVO.setAnalyseQuestion(analysePackage.getAnalyseQuestion());
            resultVO.setAnalyseWrong(analysePackage.getAnalyseWrong());
            resultVO.setSuggestLabels(analysePackage.getSuggestLabels());
            resultVO.setRecommendTopics(analysePackage.getRecommendTopics());
            
            return resultVO;
            
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            log.error("AI分析请求被中断", e);
            throw new RuntimeException("AI分析请求被中断");
        } catch (TimeoutException e) {
            log.error("AI分析超时 - RequestId: {}", baseRequest.getId(), e);
            throw new RuntimeException("AI分析超时,请稍后重试");
        } catch (ExecutionException e) {
            log.error("AI分析执行异常", e);
            throw new RuntimeException("AI分析失败: " + e.getCause().getMessage());
        } catch (Exception e) {
            log.error("发送图片URL到RabbitMQ失败", e);
            throw new RuntimeException("AI分析请求发送失败: " + e.getMessage());
        }
    }

三.功能测试

与AI端的负责开发人员进行系统联调,AI端启动,我使用Swagger发送数据,可以看到在RabbitMQ中,已经有了相关数据:

AI端也给出了回应:

总结:

本周修改了一些历史遗留的问题,并且对于错题的AI模块进行了更加深入的开发,测试了与消息队列和AI端的联调,并取得了良好的效果,这位之后的进一步开发打好了基础。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐