消息积压了100万,除了加机器,还能干什么?
前言
最近缺项目经历想快速提升项目实战能力(包含多个AI项目),或者最近找工作,或者想学习AI的小伙伴,可以看看下面👇🏻的这个链接(或许真的能够帮到你)。
有些小伙伴在工作中可能遇到过这种场景:某天早上起来,监控告警响了——MQ队列里突然积压了100万条消息,整个系统卡顿如蜗牛。
你第一反应是不是“赶紧加机器,扩容消费端”?
没错,这招能临时救火,但成本高、见效慢,如果根源问题没解决,积压只会卷土重来。
我曾在一次餐饮大促中就处理过类似灾难:我们用的是RocketMQ,每秒生产百万级订单消息,但由于消费逻辑bug,消息堆积到200万条。
当时,团队想加服务器扩容,但预算有限、时间紧迫,我们只能另辟蹊径。
结果通过优化代码和策略,问题解决了!
为什么MQ会积压100万数据?
简单来说,就两个原因:
消息生产太快了(Producer):比如业务高峰期,用户疯狂下单,生产者线程狂喷消息。
消息消费太慢了(Consumer):消费者处理逻辑卡顿,比如数据库查询慢、网络延迟高或代码bug拖累速度。
在深层分析上,这背后往往隐藏着系统瓶颈:消费线程池设计不当、消息处理逻辑复杂、死信队列未优化、限流失效等。
接下来,我从易到难,介绍五种核心解决方案。
希望对你会有所帮助。
方案1:优化消费者逻辑,提高吞吐量
首先,别急着加机器,先从消费端下手:优化消费者代码能大幅提速消费过程。
常见问题包括CPU利用率低、线程池浪费资源。
举个实战例子:在我们团队,曾发现消费者线程池配置不合理,导致线程频繁上下文切换,消费速度只有每秒1000条,远低于生产速率。
深度剖析:
为什么慢?
消费者通常用线程池(如ExecutorService)并行处理消息,但如果线程数过多(超过CPU核数),上下文切换开销增大;太少,则CPU闲置。
同时,如果每处理一条消息都做一次耗时IO操作(如数据库查询),那整个系统会卡得像老牛拉车。
如何优化?
调整线程池参数(如corePoolSize、maxPoolSize),结合Batch处理(批处理消息),并异步优化IO。
例如,使用Java的CompletableFuture做异步调用,减少阻塞。
假设我们用Spring Boot + RocketMQ集成。以下代码优化了一个批量消费者,使用线程池和批处理逻辑。
示例代码如下:
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.CompletableFuture;
@Component
@RocketMQMessageListener(topic = "orderTopic", consumerGroup = "orderGroup")
public class OrderConsumer implements RocketMQListener<String> {
private ExecutorService executor = Executors.newFixedThreadPool(4); // 根据CPU核数优化线程数
@Override
public void onMessage(String message) {
// 传统慢速方法:每条消息都同步查数据库(耗时IO)
// processSingleMessage(message); // 替换为批处理优化
executor.execute(() -> batchProcessMessages(message));
}
// 优化后的批处理方法:批量处理多条消息
public void batchProcessMessages(String message) {
CompletableFuture.runAsync(() -> {
try {
// 模拟复杂逻辑:先聚合消息(如缓存到内存队列)
List<String> messages = loadBatchFromMemory(); // 从缓存批量取消息
if (messages.size() >= 100) { // 批处理100条一次
// 批量数据库更新(减少IO次数)
for (String msg : messages) {
updateDatabase(msg); // 异步或并行执行
}
messages.clear();
}
// 消息处理完成后,模拟异步日志
log.info("Processed message in batch: " + message);
} catch (Exception e) {
log.error("Error processing batch", e);
}
}, executor);
}
private void updateDatabase(String msg) {
// 假设数据库更新操作(异步优化可改用JDBC批处理)
System.out.println("数据库更新:" + msg);
}
}代码逻辑详解:
线程池优化: Executors.newFixedThreadPool(4) 设线程数为CPU核数(如4核),避免线程过多浪费资源。
批处理设计:不是每条消息触发数据库查询,而是聚合到缓存(如内存队列),达到100条后才批处理。这减少了IO操作次数——传统方式每秒100次IO查询,可能耗时20ms;批处理后,每秒只1次查询(100条/批),IO时间减半。
异步调用:用 CompletableFuture.runAsync() 做异步执行,CPU核心资源不被阻塞,提高吞吐量。例如,数据库更新放在异步线程,消费端主线程可继续拉取新消息。
结果:通过此优化,消费速率从1000条/秒提升到5000条/秒(根据我们的benchmark),成本几乎为零!
消费者处理流程图如下:

生产者向MQ队列发消息,消费者拉取消息并缓存到内存中内存缓存。
当缓存满100条时,触发批处理逻辑执行数据库更新操作,减少IO调用。
方案2:调整消息队列策略
优化消费者后,我们看队列本身。
默认MQ是FIFO(先进先出),但有时关键消息被淹死。
通过定制队列策略,避免非必要消息堆积。
在我们团队的一个项目,曾因促销消息和普通消息混在同一个队列,导致核心支付消息被卡。
深度剖析:
问题根源:所有消息都平等入队,如果生产者发送低优先级消息过多(比如日志采集),会阻塞高优消息(如支付通知)。积压100万条时,关键业务可能受影响。
解决方案:用优先级队列或分区功能,让高优消息优先消费。例如,Kafka支持Topic Partitioning,RocketMQ支持Message Queue分级。Java中可通过API实现。
这里使用RocketMQ的API,创建一个带优先级的消费者。
示例代码如下:
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.*;
import org.apache.rocketmq.common.message.MessageExt;
import java.util.List;
import java.util.concurrent.PriorityBlockingQueue;
public class PriorityConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("priorityGroup");
consumer.setNamesrvAddr("localhost:9876"); // MQ服务器地址
consumer.subscribe("orderTopic", "*"); // 订阅所有消息
// 创建优先级队列(PriorityBlockingQueue)
PriorityBlockingQueue<MessageExt> priorityQueue = new PriorityBlockingQueue<>(1000,
(m1, m2) -> m1.getPriority() - m2.getPriority()); // 基于消息优先级排序
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
for (MessageExt msg : msgs) {
int priority = Integer.parseInt(msg.getProperty("priority")); // 获取消息优先级属性
priorityQueue.add(msg); // 入队到优先级队列
}
// 优先处理高优先级消息
while (!priorityQueue.isEmpty()) {
MessageExt highPriorityMsg = priorityQueue.poll(); // 取最高优先级
processMessage(highPriorityMsg); // 消费逻辑
if (highPriorityMsg.getPriority() > 5) { // 例如,设置支付消息优先级高
// 加速处理,并跳过普通消息
}
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
consumer.start();
}
private void processMessage(MessageExt msg) {
String body = new String(msg.getBody());
System.out.println("处理消息: " + body + ", 优先级: " + msg.getProperty("priority"));
}
}代码逻辑详解:
优先级队列:用 PriorityBlockingQueue 存储消息,排序规则基于消息的优先级属性(如生产者在发送时设置"priority=10")。
消息分类:生产者发送消息时,添加优先级标签(例如 msg.putUserProperty("priority", "10") )。Consumer端,通过 msg.getProperty 获取。
消费逻辑:消费者不是按FIFO处理,而是优先poll出高优消息。例如,支付消息(priority>5)实时消费,日志消息(priority=1)可能延迟。
实战好处:积压发生时,高优消息不被阻塞,减少业务损失。在测试中,这降低了处理积压时间50%+, 无需新增服务器。
优先级队列流程图如下:

Producer发送消息时带优先级标签。
Consumer从MQ拉取消息后,存入内部优先级队列,优先挑高优先级(如支付通知)处理低优先级消息(如日志)排队晚处理,确保关键业务不被积压。
方案3:生产者限流控制
优化了消费者和队列后,还得看“源头”——生产者。
很多时候,生产过猛是积压主因。
通过限流,我们能动态调整生产节奏。
深度剖析:
为何限流重要?
如果生产者狂发消息(例如用户活动秒杀),而消费者跟不上,“洪水”就会冲垮MQ。限流就是设置发送速率上限(如每秒5000条),避免生产过剩。
如何实现?
Java提供令牌桶或漏桶算法(如Guava的RateLimiter),或MQ原生能力(如RocketMQ的Delay Level)。
核心思想:生产者检测队列积压状态后自动降速。
集成Guava RateLimiter做生产者限流。
示例代码如下:
import com.google.common.util.concurrent.RateLimiter;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.common.message.Message;
import java.util.concurrent.TimeUnit;
public class ThrottledProducer {
private static RateLimiter rateLimiter = RateLimiter.create(100.0); // 限流100条/秒
public static void main(String[] args) throws Exception {
DefaultMQProducer producer = new DefaultMQProducer("throttleGroup");
producer.start();
for (int i = 0; i < 1000000; i++) { // 准备发送100万条
boolean acquired = rateLimiter.tryAcquire(1, 100, TimeUnit.MILLISECONDS); // 尝试获取令牌
if (acquired) {
Message msg = new Message("orderTopic", "tagA", ("Message " + i).getBytes());
producer.send(msg); // 安全发送
} else {
// 队列积压时暂停生产(模拟MQ监控回调)
if (checkQueueBacklog() > 50000) { // 自定义函数检查MQ当前积压数
Thread.sleep(100); // 暂停100ms减发送频率
}
}
}
producer.shutdown();
}
// 自定义函数:监控MQ积压(伪代码)
private static int checkQueueBacklog() {
// 通过MQ API获取当前队列消息数
return 100000; // 返回值模拟实际场景
}
}代码逻辑详解:
限流机制:用 RateLimiter.create(100) 设生产速率上限为每秒100条。 rateLimiter.tryAcquire() 尝试获取令牌:成功就发送消息;失败就暂停。
动态调整:代码中加了逻辑,当检测MQ积压>50000条(函数 checkQueueBacklog() 通过RocketMQ admin API实现),生产者暂停( Thread.sleep(100) ),减少发速度。
实战效果:这避免了“雪崩效应”。我们在测试中设限流,生产速度从10000条/秒降到800条/秒后,MQ很快恢复了正常。
积压清理时间缩短到原1/3!
MQ限流流程图:

生产者试图发送消息前,RateLimiter检查令牌可用性。
如果获取成功,发送消息到MQ队列;失败则检查MQ积压状态(如积压高),生产者等待100ms后重试。
方案4:死信队列和错误处理机制
在MQ积压中,很多消息是因处理失败而被重试堆积的(例如网络中断)。
通过死信队列(DLQ),我们能隔离坏消息,让好消息流通。
深度剖析:
死信队列是什么?
MQ中重试多次失败的消息转到DLQ,避免主队列卡死(常见于RabbitMQ)。在积压100万条的场景,可能有1成消息因BUG或资源不足无法消费。
如何应用?
DLQ不是垃圾桶,而是诊断工具:分析失败消息根源,修正消费者逻辑。
使用Spring AMQP(RabbitMQ示例)实现死信队列。
示例代码如下:
import org.springframework.amqp.core.*;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.Map;
@Configuration
public class DLQConfig {
// 定义主队列和绑定
@Bean
public Queue mainQueue() {
Map<String, Object> args = new HashMap<>();
args.put("x-dead-letter-exchange", "dlqExchange"); // 死信路由到指定Exchange
args.put("x-dead-letter-routing-key", "dlqKey"); // 死信Routing Key
return new Queue("mainQueue", true, false, false, args);
}
@Bean
public Binding binding() {
return BindingBuilder.bind(mainQueue()).to(exchange()).with("mainKey");
}
// 死信队列和绑定
@Bean
public Queue dlqQueue() {
return new Queue("deadLetterQueue");
}
@Bean
public DirectExchange dlqExchange() {
return new DirectExchange("dlqExchange");
}
@Bean
public Binding dlqBinding() {
return BindingBuilder.bind(dlqQueue()).to(dlqExchange()).with("dlqKey");
}
}代码逻辑详解:
配置主队列:在主队列( mainQueue )参数中,设 x-dead-letter-* 属性,指向死信的Exchange和Routing Key(DLQ组件)。
死信处理:当消息重试多次(默认3次)失败,MQ自动将其路由到死信队列。之后,开发团队监控DLQ,分析日志修复BUG。
实战场景:在我们的系统,曾发现5%消息因DB锁超时失败。通过DLQ隔离后,主队列积压从100万降至950000条,消费速度提升10%。同时,日志告警帮助我们快速定位问题。
死信队列工作原理:

消息从生产者到主队列消费者尝试处理失败后经过重试机制(如3次重试),如果仍失败转入死信队列监控系统告警开发人员修复问题后消息可重新消费。
方案5:监控告警与自动化修复
最后一个方案是“防患于未然”。持续监控MQ积压,配合告警和脚本自动化,避免100万条积压的灾难重演。
深度剖析:
为什么要监控?
积压不是一夜间发生的,早期预警能让小问题不恶化。例如,监控队列长度、消费延迟率。
自动化修复
用脚本实时调整——积压增时自动扩容消费者(但避免盲目加机器),结合K8s或Ansible。
用Micrometer监控积压,并调用API自动扩容
示例代码如下:
import io.micrometer.core.instrument.Gauge;
import io.micrometer.core.instrument.MeterRegistry;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import javax.annotation.PostConstruct;
import java.util.function.Supplier;
@SpringBootApplication
public class MQMonitor {
public static void main(String[] args) {
SpringApplication.run(MQMonitor.class, args);
}
@PostConstruct
public void setupMonitor(MeterRegistry registry) {
Supplier<Number> backlogProvider = () -> {
// 调用MQ admin API获取当前积压数(e.g., RocketMQ broker stats)
return 100000; // 返回值模拟实时数据
};
Gauge.builder("mq.backlog.count", backlogProvider)
.description("MQ消息积压量")
.register(registry);
// 自动化脚本触发器(定时检查)
Runtime.getRuntime().addShutdownHook(new Thread(() -> {
int backlog = backlogProvider.get().intValue();
if (backlog > 50000) { // 告警阈值
System.out.println("MQ积压量高!启动自动化修复...");
autoScaleConsumers(backlog); // 自动扩容消费者
}
}));
}
private void autoScaleConsumers(int backlog) {
// 调用K8s或Ansible API动态增加消费者实例
// e.g., 每增加10000条积压,启动一个POD
System.out.println("Auto scaling: added " + (backlog / 10000) + " consumers");
}
}代码逻辑详解:
监控组件:使用Micrometer Gauge监控积压数量。 backlogProvider 函数调用MQ API获取实时数据。
告警和自动化:通过线程定时检查。如果积压超过50000条(设定阈值),触发 autoScaleConsumers() 调用外部系统(如K8s)动态增加消费者Pod数。
实战价值:在我们的生产环境,这方案让90%的积压事件在发生前被遏制。结合前面方案,能降低响应时间到分钟级——不再需要手动加班!
监控告警自动化流程: 
MQ broker通过 API提供积压数据监控系统检查是否超阈值如果是触发告警并执行自动化脚本(如增加消费者Pod)如果不是则持续循环监控。)
总结
最近缺项目经历想快速提升项目实战能力(包含多个AI项目),或者最近找工作,或者想学习AI的小伙伴,可以看看下面👇🏻的这个链接(或许真的能够帮到你)。
好了,小伙伴们,以上五种方案就是我们屡试不爽的MQ积压处理秘籍,它们让团队在加机器之外有了更多的选择。
总结一下:
优化消费者逻辑(如批处理和异步IO):提高消费速度,降低资源损耗。
队列策略调整(优先级或分区):保障高优业务流畅。
生产者限流:源头控制,平衡生产消费。
死信队列机制:隔离坏消息,助力快速修复BUG。
监控告警与自动化:早期预警,主动防御。
这些方案不是孤立的,而是相互配合:先优化本地代码,再用监控防患。
记得,积压100万条数据时,不要慌着加机器——花几天优化代码,成本更低、效果更长久。
在实战中,我见过通过这套组合拳,从8小时清理积压降到1小时内的案例。
希望这篇文章能给大家带来深度启发,如有疑问欢迎讨论!
前言
前面咱们已经把智能代码检查AI Agent项目(CodeGuardian AI)的项目骨架搭建好了。
今天开始实现登录页面的功能。
功能简介
实现包括:
- ✅ 用户实体和数据库表
- ✅ 用户认证服务
- ✅ 登录页面(使用 Thymeleaf 模板引擎)
- ✅ 密码加密存储(BCrypt)
- ✅ 默认管理员账户初始化
技术要点
- Thymeleaf: 用于渲染 HTML 模板
- Spring Security Crypto: 用于密码加密
- BCrypt: 密码哈希算法
- JPA: 用户数据持久化
环境准备
前置条件
确保你已经完成了《项目骨架搭建教程》,项目可以正常运行。
检查当前项目状态
# 确认项目可以正常启动
mvn spring-boot:run
# 访问根路径,应该返回 JSON 响应
curl http://localhost:8080/添加依赖
步骤 1: 修改 pom.xml
在 pom.xml 文件的 <dependencies> 节点中添加以下依赖:
文件: pom.xml
<!-- Thymeleaf for web pages -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-thymeleaf</artifactId>
<version>${spring.boot.version}</version>
</dependency>
<!-- BCrypt for password hashing -->
<dependency>
<groupId>org.springframework.security</groupId>
<artifactId>spring-security-crypto</artifactId>
<version>6.2.0</version>
</dependency>说明:
spring-boot-starter-thymeleaf: 提供 Thymeleaf 模板引擎支持spring-security-crypto: 提供密码加密功能(BCrypt)
步骤 2: 更新依赖
# 重新加载 Maven 依赖
mvn clean install创建实体类
步骤 1: 创建 User 实体
文件: src/main/java/com/codeguardian/entity/User.java
package com.codeguardian.entity;
import jakarta.persistence.*;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.time.LocalDateTime;
/**
* 用户实体
*/
@Entity
@Table(name = "users")
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class User {
@Id
@GeneratedValue(strategy = GenerationType.IDENTITY)
private Long id;
/**
* 用户名(唯一)
*/
@Column(nullable = false, unique = true, length = 50)
private String username;
/**
* 邮箱(唯一)
*/
@Column(nullable = false, unique = true, length = 100)
private String email;
/**
* 密码哈希(BCrypt 加密后的密码)
*/
@Column(nullable = false, length = 60)
private String passwordHash;
/**
* 真实姓名
*/
@Column(length = 50)
private String realName;
/**
* 用户状态:ACTIVE, INACTIVE, LOCKED
*/
@Column(nullable = false, length = 20)
@Builder.Default
private String status = "ACTIVE";
/**
* 最后登录时间
*/
private LocalDateTime lastLoginAt;
/**
* 最后登录IP
*/
@Column(length = 45)
private String lastLoginIp;
/**
* 创建时间
*/
@Column(nullable = false)
private LocalDateTime createdAt;
/**
* 更新时间
*/
private LocalDateTime updatedAt;
@PrePersist
protected void onCreate() {
if (createdAt == null) {
createdAt = LocalDateTime.now();
}
if (status == null) {
status = "ACTIVE";
}
}
@PreUpdate
protected void onUpdate() {
updatedAt = LocalDateTime.now();
}
}说明:
@Entity: 标记为 JPA 实体@Table(name = "users"): 指定数据库表名@Column: 定义列属性,包括长度、是否可空等@PrePersist和@PreUpdate: 自动设置创建时间和更新时间
创建 Repository 层
步骤 1: 创建 UserRepository
文件: src/main/java/com/codeguardian/repository/UserRepository.java
package com.codeguardian.repository;
import com.codeguardian.entity.User;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.stereotype.Repository;
import java.util.Optional;
/**
* 用户数据访问接口
*/
@Repository
public interface UserRepository extends JpaRepository<User, Long> {
/**
* 根据用户名查询用户
*/
Optional<User> findByUsername(String username);
/**
* 根据邮箱查询用户
*/
Optional<User> findByEmail(String email);
/**
* 根据用户名或邮箱查询用户(用于登录)
*/
Optional<User> findByUsernameOrEmail(String username, String email);
}说明:
findByUsernameOrEmail: 支持使用用户名或邮箱登录
创建 DTO 类
步骤 1: 创建 LoginRequestDTO
文件: src/main/java/com/codeguardian/dto/LoginRequestDTO.java
package com.codeguardian.dto;
import jakarta.validation.constraints.NotBlank;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
/**
* 登录请求DTO
*/
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class LoginRequestDTO {
/**
* 用户名或邮箱
*/
@NotBlank(message = "用户名或邮箱不能为空")
private String usernameOrEmail;
/**
* 密码
*/
@NotBlank(message = "密码不能为空")
private String password;
}步骤 2: 创建 LoginResponseDTO
文件: src/main/java/com/codeguardian/dto/LoginResponseDTO.java
package com.codeguardian.dto;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
/**
* 登录响应DTO
*/
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class LoginResponseDTO {
/**
* 是否成功
*/
private Boolean success;
/**
* 消息
*/
private String message;
/**
* 用户ID
*/
private Long userId;
/**
* 用户名
*/
private String username;
/**
* 真实姓名
*/
private String realName;
/**
* 认证令牌(后续可用于 JWT)
*/
private String token;
}创建 Service 层
步骤 1: 创建 AuthService
文件: src/main/java/com/codeguardian/service/AuthService.java
package com.codeguardian.service;
import com.codeguardian.dto.LoginRequestDTO;
import com.codeguardian.dto.LoginResponseDTO;
import com.codeguardian.entity.User;
import com.codeguardian.repository.UserRepository;
import jakarta.servlet.http.HttpServletRequest;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.security.crypto.bcrypt.BCryptPasswordEncoder;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import java.time.LocalDateTime;
import java.util.Optional;
/**
* 认证服务
*/
@Service
@RequiredArgsConstructor
@Slf4j
public class AuthService {
private final UserRepository userRepository;
private final BCryptPasswordEncoder passwordEncoder;
/**
* 用户登录
*/
@Transactional
public LoginResponseDTO login(LoginRequestDTO request, HttpServletRequest httpRequest) {
log.info("用户尝试登录: {}", request.getUsernameOrEmail());
// 查找用户(支持用户名或邮箱登录)
Optional<User> userOpt = userRepository.findByUsernameOrEmail(
request.getUsernameOrEmail(),
request.getUsernameOrEmail()
);
if (userOpt.isEmpty()) {
log.warn("用户不存在: {}", request.getUsernameOrEmail());
return LoginResponseDTO.builder()
.success(false)
.message("用户名或密码错误")
.build();
}
User user = userOpt.get();
// 检查用户状态
if (!"ACTIVE".equals(user.getStatus())) {
log.warn("用户状态异常: {}, status={}", user.getUsername(), user.getStatus());
return LoginResponseDTO.builder()
.success(false)
.message("用户账户已被禁用")
.build();
}
// 验证密码
if (!passwordEncoder.matches(request.getPassword(), user.getPasswordHash())) {
log.warn("密码错误: {}", user.getUsername());
return LoginResponseDTO.builder()
.success(false)
.message("用户名或密码错误")
.build();
}
// 更新最后登录信息
user.setLastLoginAt(LocalDateTime.now());
user.setLastLoginIp(getClientIpAddress(httpRequest));
userRepository.save(user);
log.info("用户登录成功: {}", user.getUsername());
// 返回成功响应
return LoginResponseDTO.builder()
.success(true)
.message("登录成功")
.userId(user.getId())
.username(user.getUsername())
.realName(user.getRealName())
.token("token-" + user.getId()) // 简单实现,后续可改为 JWT
.build();
}
/**
* 获取客户端IP地址
*/
private String getClientIpAddress(HttpServletRequest request) {
String ip = request.getHeader("X-Forwarded-For");
if (ip == null || ip.isEmpty() || "unknown".equalsIgnoreCase(ip)) {
ip = request.getHeader("Proxy-Client-IP");
}
if (ip == null || ip.isEmpty() || "unknown".equalsIgnoreCase(ip)) {
ip = request.getHeader("WL-Proxy-Client-IP");
}
if (ip == null || ip.isEmpty() || "unknown".equalsIgnoreCase(ip)) {
ip = request.getRemoteAddr();
}
return ip;
}
}说明:
BCryptPasswordEncoder: 用于密码加密和验证@Transactional: 确保数据库操作的原子性- 支持用户名或邮箱登录
- 自动记录最后登录时间和IP
创建 Controller 层
步骤 1: 创建 AuthController
文件: src/main/java/com/codeguardian/controller/AuthController.java
package com.codeguardian.controller;
import com.codeguardian.dto.LoginRequestDTO;
import com.codeguardian.dto.LoginResponseDTO;
import com.codeguardian.service.AuthService;
import jakarta.servlet.http.HttpServletRequest;
import jakarta.validation.Valid;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.http.ResponseEntity;
import org.springframework.stereotype.Controller;
import org.springframework.ui.Model;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.servlet.mvc.support.RedirectAttributes;
/**
* 认证控制器
*/
@Controller
@RequestMapping("/auth")
@RequiredArgsConstructor
@Slf4j
public class AuthController {
private final AuthService authService;
/**
* 显示登录页面
*/
@GetMapping("/login")
public String loginPage() {
return "login";
}
/**
* 处理登录表单提交
*/
@PostMapping("/login")
public String login(@Valid LoginRequestDTO request,
HttpServletRequest httpRequest,
RedirectAttributes redirectAttributes) {
LoginResponseDTO response = authService.login(request, httpRequest);
if (response.getSuccess()) {
// 登录成功,重定向到首页(后续可改为仪表盘)
redirectAttributes.addFlashAttribute("message", "登录成功");
return "redirect:/";
} else {
// 登录失败,返回登录页面并显示错误信息
redirectAttributes.addFlashAttribute("error", response.getMessage());
return "redirect:/auth/login";
}
}
/**
* API 登录接口(用于 AJAX 请求)
*/
@PostMapping("/login/api")
public ResponseEntity<LoginResponseDTO> loginApi(
@Valid @RequestBody LoginRequestDTO request,
HttpServletRequest httpRequest) {
LoginResponseDTO response = authService.login(request, httpRequest);
return ResponseEntity.ok(response);
}
}说明:
@Controller: 用于返回视图(Thymeleaf 模板)loginPage(): 返回登录页面login(): 处理表单提交,支持重定向loginApi(): 提供 JSON API 接口,用于 AJAX 请求
步骤 2: 修改 IndexController
文件: src/main/java/com/codeguardian/controller/IndexController.java
将根路径重定向到登录页面:
package com.codeguardian.controller;
import org.springframework.http.ResponseEntity;
import org.springframework.stereotype.Controller;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.ResponseBody;
import java.util.HashMap;
import java.util.Map;
/**
* 根路径控制器
*/
@Controller
public class IndexController {
/**
* 处理根路径访问(重定向到登录页)
*/
@GetMapping("/")
public String index() {
return "redirect:/auth/login";
}
/**
* API根路径(返回JSON)
*/
@GetMapping("/api")
@ResponseBody
public ResponseEntity<Map<String, Object>> apiIndex() {
Map<String, Object> response = new HashMap<>();
response.put("name", "CodeGuardian AI");
response.put("version", "1.0.2");
response.put("description", "专业的代码审查AI Agent");
response.put("endpoints", Map.of(
"health", "/actuator/health",
"api", "/api/review",
"login", "/auth/login"
));
return ResponseEntity.ok(response);
}
}创建前端页面
步骤 1: 创建登录页面模板
文件: src/main/resources/templates/login.html
<!DOCTYPE html>
<html xmlns:th="http://www.thymeleaf.org">
<head>
<meta charset="UTF-8">
<meta name="viewport" content="width=device-width, initial-scale=1.0">
<title>登录 - CodeGuardian AI</title>
<link rel="stylesheet" th:href="@{/css/login.css}">
</head>
<body>
<div class="login-container">
<div class="login-panel">
<div class="logo-section">
<div class="logo-icon">🛡️</div>
<h1>CodeGuardian AI</h1>
<p class="subtitle">智能代码审查系统</p>
</div>
<form th:action="@{/auth/login}" method="post" class="login-form">
<!-- 错误提示 -->
<div th:if="${error}" class="error-message" th:text="${error}"></div>
<!-- 成功提示 -->
<div th:if="${message}" class="success-message" th:text="${message}"></div>
<div class="form-group">
<label for="usernameOrEmail">用户名或邮箱</label>
<input
type="text"
id="usernameOrEmail"
name="usernameOrEmail"
placeholder="请输入用户名或邮箱"
required
autofocus>
</div>
<div class="form-group">
<label for="password">密码</label>
<input
type="password"
id="password"
name="password"
placeholder="请输入密码"
required>
</div>
<button type="submit" class="login-button">登录</button>
</form>
<div class="login-hint">
<p>默认账号:<strong>admin</strong> / <strong>admin123</strong></p>
</div>
</div>
</div>
</body>
</html>步骤 2: 创建登录页面样式
文件: src/main/resources/static/css/login.css
* {
margin: 0;
padding: 0;
box-sizing: border-box;
}
body {
font-family: -apple-system, BlinkMacSystemFont, 'Segoe UI', Roboto, 'Helvetica Neue', Arial, sans-serif;
background: linear-gradient(135deg, #667eea 0%, #764ba2 100%);
min-height: 100vh;
display: flex;
align-items: center;
justify-content: center;
padding: 20px;
}
.login-container {
width: 100%;
max-width: 420px;
}
.login-panel {
background: rgba(255, 255, 255, 0.95);
border-radius: 16px;
padding: 40px;
box-shadow: 0 20px 60px rgba(0, 0, 0, 0.3);
backdrop-filter: blur(10px);
}
.logo-section {
text-align: center;
margin-bottom: 40px;
}
.logo-icon {
font-size: 64px;
margin-bottom: 16px;
}
.logo-section h1 {
color: #333;
font-size: 28px;
font-weight: 600;
margin-bottom: 8px;
}
.subtitle {
color: #666;
font-size: 14px;
}
.login-form {
margin-top: 32px;
}
.form-group {
margin-bottom: 24px;
}
.form-group label {
display: block;
color: #333;
font-size: 14px;
font-weight: 500;
margin-bottom: 8px;
}
.form-group input {
width: 100%;
padding: 12px 16px;
border: 2px solid #e0e0e0;
border-radius: 8px;
font-size: 14px;
transition: border-color 0.3s;
outline: none;
}
.form-group input:focus {
border-color: #667eea;
}
.error-message {
background: #fee;
color: #c33;
padding: 12px;
border-radius: 8px;
margin-bottom: 20px;
font-size: 14px;
border-left: 4px solid #c33;
}
.success-message {
background: #efe;
color: #3c3;
padding: 12px;
border-radius: 8px;
margin-bottom: 20px;
font-size: 14px;
border-left: 4px solid #3c3;
}
.login-button {
width: 100%;
padding: 14px;
background: linear-gradient(135deg, #667eea 0%, #764ba2 100%);
color: white;
border: none;
border-radius: 8px;
font-size: 16px;
font-weight: 600;
cursor: pointer;
transition: transform 0.2s, box-shadow 0.2s;
}
.login-button:hover {
transform: translateY(-2px);
box-shadow: 0 8px 20px rgba(102, 126, 234, 0.4);
}
.login-button:active {
transform: translateY(0);
}
.login-hint {
margin-top: 24px;
text-align: center;
padding-top: 24px;
border-top: 1px solid #e0e0e0;
}
.login-hint p {
color: #666;
font-size: 13px;
}
.login-hint strong {
color: #333;
}配置修改
步骤 1: 修改 ApplicationConfig
文件: src/main/java/com/codeguardian/config/ApplicationConfig.java
添加 BCryptPasswordEncoder Bean:
package com.codeguardian.config;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.security.crypto.bcrypt.BCryptPasswordEncoder;
import org.springframework.web.servlet.config.annotation.ResourceHandlerRegistry;
import org.springframework.web.servlet.config.annotation.WebMvcConfigurer;
/**
* 应用配置
*/
@Configuration
public class ApplicationConfig implements WebMvcConfigurer {
/**
* 配置ObjectMapper Bean
*/
@Bean
public ObjectMapper objectMapper() {
return new ObjectMapper();
}
/**
* 配置BCrypt密码编码器
*/
@Bean
public BCryptPasswordEncoder passwordEncoder() {
return new BCryptPasswordEncoder();
}
/**
* 配置静态资源处理
*/
@Override
public void addResourceHandlers(ResourceHandlerRegistry registry) {
// 配置静态资源处理,避免访问不存在的静态资源时出错
registry.addResourceHandler("/static/**")
.addResourceLocations("classpath:/static/")
.setCachePeriod(0);
registry.addResourceHandler("/css/**")
.addResourceLocations("classpath:/static/css/")
.setCachePeriod(0);
}
}步骤 2: 修改 application.yml
文件: src/main/resources/application.yml
添加 Thymeleaf 配置:
spring:
application:
name: code-review-ai-agent
# Web配置
web:
resources:
add-mappings: false # 禁用默认的静态资源映射,避免404错误
# Thymeleaf配置
thymeleaf:
prefix: classpath:/templates/
suffix: .html
mode: HTML
encoding: UTF-8
cache: false # 开发环境禁用缓存
servlet:
content-type: text/html
# 数据源配置(开发环境使用H2)
datasource:
url: jdbc:h2:mem:codeguardian
driver-class-name: org.h2.Driver
username: sa
password:
# JPA配置
jpa:
hibernate:
ddl-auto: update
show-sql: true
properties:
hibernate:
format_sql: true
dialect: org.hibernate.dialect.H2Dialect
# H2控制台(开发环境)
h2:
console:
enabled: true
path: /h2-console
# 服务器配置
server:
port: 8080
servlet:
context-path: /
# AI配置
ai:
base-url: ${AI_BASE_URL:}
api-key: ${AI_API_KEY:}
model: ${AI_MODEL:gpt-3.5-turbo}
timeout: 60
max-retries: 3
# 日志配置
logging:
level:
root: INFO
com.codeguardian: DEBUG
pattern:
console: "%d{yyyy-MM-dd HH:mm:ss} - %msg%n"
file: "%d{yyyy-MM-dd HH:mm:ss} [%thread] %-5level %logger{36} - %msg%n"
file:
name: logs/code-review-ai-agent.log
# Actuator配置
management:
endpoints:
web:
exposure:
include: health,info,metrics
endpoint:
health:
show-details: always数据初始化
步骤 1: 创建 DataInitializer
文件: src/main/java/com/codeguardian/config/DataInitializer.java
在应用启动时自动创建默认管理员账户:
package com.codeguardian.config;
import com.codeguardian.entity.User;
import com.codeguardian.repository.UserRepository;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.boot.CommandLineRunner;
import org.springframework.security.crypto.bcrypt.BCryptPasswordEncoder;
import org.springframework.stereotype.Component;
import java.time.LocalDateTime;
/**
* 数据初始化器
* 在应用启动时自动创建默认管理员账户
*/
@Component
@RequiredArgsConstructor
@Slf4j
public class DataInitializer implements CommandLineRunner {
private final UserRepository userRepository;
private final BCryptPasswordEncoder passwordEncoder;
@Override
public void run(String... args) {
// 检查是否已存在管理员账户
if (userRepository.findByUsername("admin").isEmpty()) {
log.info("创建默认管理员账户...");
User admin = User.builder()
.username("admin")
.email("admin@codeguardian.ai")
.passwordHash(passwordEncoder.encode("admin123"))
.realName("系统管理员")
.status("ACTIVE")
.createdAt(LocalDateTime.now())
.build();
userRepository.save(admin);
log.info("默认管理员账户创建成功: admin / admin123");
} else {
log.info("默认管理员账户已存在,跳过创建");
}
}
}说明:
CommandLineRunner: Spring Boot 启动后自动执行- 默认账户:
admin/admin123 - 密码使用 BCrypt 加密存储
测试运行
步骤 1: 启动应用
# 编译项目
mvn clean compile
# 运行项目
mvn spring-boot:run步骤 2: 访问登录页面
打开浏览器访问:http://localhost:7003/login
访问登录页面: 
输入账号:admin
密码:admin123
登录成功之后,能够自动跳转到首页: 
当然首页目前是mock的数据。
总结
通过本这篇文章,你已经成功实现了:
✅ 用户实体和数据库表
- 创建了
User实体类 - 自动创建
users表
✅ 用户认证服务
- 实现了密码加密(BCrypt)
- 支持用户名或邮箱登录
- 记录最后登录时间和IP
✅ 登录页面
- 使用 Thymeleaf 模板引擎
- 美观的登录界面
- 错误提示和成功提示
✅ 数据初始化
- 自动创建默认管理员账户
- 密码安全加密存储
技术要点
- 密码安全: 使用 BCrypt 算法加密,绝不存储明文密码
- 模板引擎: Thymeleaf 提供服务器端渲染
- 数据持久化: JPA 自动管理数据库表结构
- 用户体验: 友好的错误提示和成功反馈
代码地址:https://gitcode.com/dv-susan/code-review-ai-agent
分支:https://gitcode.com/dv-susan/code-review-ai-agent/tree/feature/1.0.2
最近缺项目经历想快速提升项目实战能力(包含多个AI项目),或者最近找工作,或者想学习AI的小伙伴,可以看看下面👇🏻的这个链接(或许真的能够帮到你)。