考试系统交接日记(五):高并发架构简述与kafka处理高并发业务
一、高并发架构简述
考试系统是一个百万级高并发的系统,可能同时有百万级的用户参加考试。
考试系统主要涉及到高并发的功能为:用户报名考试、获取试卷、用户提交试卷答案、返回用户得分。
考试系统其余主要功能为:管理员与老师管理试卷、老师查看自己班的考生得分、考试记录推送到其它系统等。
考试系统使用了30多个服务器,其中某些服务器上(性能比较好的),1个服务器部署2个项目实例。
考试系统后台大部分为springboot项目。
●架构记录如下:
1.使用nginx做网关,实现负载均衡、静态资源访问、后台请求转发。
使用2台服务器,2个nginx;1个供老师、学员、管理员入口使用,1个供IM及时通讯使用。
2.前端项目中,将index.html等项目放到nginx服务器指定目录中,剩余的大部分文件上传到CDN上。
3.后端项目中,按照业务划分,分为5个子项目:
(1)基本信息项目,提供用户基本信息等,使用dubbo与其他项目通信;在7台服务器上部署了12个实例。
(2)试卷标签项目,提供给试卷或题目加标签以及查询搜索等服务,使用http(Controller)与其他项目通信;在5台服务器上部署了6个实例。
(3)学员考试项目,提供学员报名考试、获取试卷、提交答案、获取分数等服务,仅为移动端后台项目;在9台服务器上部署了14个实例。
(4)PC端管理系统项目,供老师或管理员查询自己班级内学员基本信息、考试信息、编辑试卷、导出报表等使用,并发量不高;在2台服务器上部署了2个实例,其中1个实例供基本使用,另一个实例仅供导出报表使用。
(5)移动端管理系统项目,与PC端管理项目管理系统项目类似;在2台服务器上部署了2个实例。
4.在2台服务器上搭建了Mysql服务器,主从关系。
5.在2台服务器上搭建了Redis服务器,集群关系。
6.在3台服务器上搭建了Kafka消息队列,集群关系。(旧版本Kafka,使用时必须安装zookeeper)
使用消息队列实现异步操作,缓解高并发压力。
主要在学员报名考试、学员提交考试答案、老师修改试卷导致缓存失效通知事件等。
7.在3台服务器上搭建了3个IM即时通信项目,GoLang;作用如下:
当学员提交答案时,存入kafka队列(异步操作解决高并发问题),服务器从kafka队列获得信息、计算完分数后,主动将分数回传给学员。
这时服务器调用此IM项目。(基于websocket)
8.在1台服务器上部署了springcloud-config-server项目,用于统一管理配置文件。配置文件在git上,此项目从git上获取配置文件信息,其它项目从这里获取配置文件信息;当配置文件改变时,此项目也可以主动将修改后的配置文件推送给其它项目。
9.在2台服务器上部署了2个单独项目:考试成绩推送;用来将考试成绩推送到其它系统。从Kafka接收信息、用RabbitMQ推送给其它系统。
10.高并发项目支持扩展,扩展方法如下:
(1)升级服务器硬件资源,同时修改启动时的jvm参数,增大-Xms与-Xmx。
(2)增加后台服务器数量,同时修改调用该服务的Dubbo配置、nginx网关信息,增加新IP。
(3)增加mysql数据库服务器数量,同时修改mysql配置文件,例如增加minimun-idle、minimun-pool-size等。
二、Kafka处理高并发业务
在考试系统-学员考试项目中,学员报名考试、获取试卷信息、提交考试答案等操作,都属于百万级高并发操作,如果采用同步处理方式,很可能造成链接数过多、服务器压力过大而宕机。
因此借助Kafka消息队列实现异步处理。
●以学员报名考试为例,具体如下:
1.学员在前端点击"考试报名"按钮,向后台发起请求;页面显示"请等待"转圈图标。
2.学员请求先走身份验证Controller,header传参,参数md5加密并生成sign签名(前端加密,后端按照同样方法加密,比较加密后的sign,一致则有效),服务器验证身份后,将用户信息存入session与redis(短信息在session中,长信息在redis中),然后让跳转到考试报名Controller。
3.此时,请求通过拦截器Interceptor,拦截器从session、redis中获取用户信息,存入ThreadLocal,供后续使用;然后跳转到参加考试Controller。
如果没有获取到用户信息,则说明用户无权访问该页面,直接跳转到错误页面。
4.(1)进入考试报名Controller,这里用到了使用【redis限制并发调用的自定义方法lock】,防止相同参数的请求因为某些异常同时多次请求该接口造成错误数据导致处理异常,后续进行讲解。
(2)进行简单处理,用redis存取信息、封装javabean,然后扔入Kafka队列,并返回给前端消息。
(3)此时,前端页面停止显示"请等待"转圈图标,出现"参加考试"按钮。(点击后会请求下一个Controller,这里略)
5.另一个service从Kafka获得学员报名消息,使用【CompletableFuture.runAsync()】方法,进行异步处理;之后继续监听Kafka消息队列。
6.CompletableFuture.runAsync()方法中,再进行具体的学员报名信息处理操作,分析信息,存入数据库,更新redis。
●redis限制并发调用的自定义方法lock(防止相同参数的请求因为某些异常同时多次请求该接口造成错误数据导致处理异常):
XXXServiceImpl.java中:
//样例
String examId = "1234";
String userId = "abcd";
//redis的key的格式,join是方法名的标识
String formatStr = "exam:%s:user:%s:join";
//得到一个redisKey,String类型
String redisKey = String.format(formatStr, examId, userId);
//注入一个自定义锁类
@Autowired
private SimpleDistributedLock lock;
//然后有一个加锁方法
lock.lock(redisKey);
try{
//调用某个方法,join0
return join0(examId, userId);
} finally{
//调用完后,释放锁;不管是否出现异常都会释放。
lock.releaseLock(redisKey);
}
自定义锁类SimpleDistributedLock.java:
@Component
@Slf4j
public class SimpleDistributedLock {
//redis操作类
private final StringRedisTemplate redis;
//构造方法
public SimpleDistributedLock(StringRedisTemplate redis){
this.redis = redis;
}
//加锁方法
public void lock(String key){
do {
//redis中,以key为key,以"1"为值,存1天
//如果reids中不存在这个key,则保存成功,返回true
//如果reids中已存在这个key,则不进行保存,返回false
Boolean acquired = redis.opsForValue().setIfAbsent(key,"1",1,TimeUnit.DAYS);
//正常情况下不走这个,除非redis炸了
if (acquired == null){
throw new IllegalStateException();
}
//如果保存成功,则加锁成功,跳出该方法,继续执行后续方法
if (acquired){
break;
}
//否则,说明保存失败,有其它线程正在使用这个key当锁
else{
try{
//睡一会,然后继续循环,直到其它线程释放了这个key锁
TimeUnit.MILLISECONDS.sleep(100);
}catch(Exception e){
log.error(e.getMessage(),e );
}
}
}
while (true);
}
//释放锁,也就是从redis中删除这个key
public void releaseLock(String key){
redis.delete(key);
}
//尝试加锁的方法,如果没有锁则进行上锁操作;最后返回结果
public boolean tryLock(String key){
Boolean acquired = redis.opsForValue().setIfAbsent(key,"1",1,TimeUnit.DAYS);
//正常情况下不走这个,除非redis炸了
if (acquired == null){
throw new IllegalStateException();
}
//返回true为上锁成功,返回false为上锁失败
return acquired;
}
}
●异步方法CompletableFuture.runAsync()用法:
@Service
@Slf4j
public class MyAsyncService {
private final ExecutorService joinExecutor;
public MyAsyncService(){
//初始化executor,corePoolSize为10,maximumPoolSize为10,keepAliveTimer为10,单位为毫秒
this.joinExecutor = new ThreadPoolExecutor(10, 10, 10, TimeUnit.MILLISECONDS,
new SynchronousQueue<>(),
new ThreadFactoryBuilder().setNameFormat("MyAsyncService-join-pool-%s-thread").build(),
new ThreadPoolExecutor.CallerRunsPolicy());
}
@PreDestroy
public void destroy() throws InterruptedException {
joinExecutor.shutdown();
//timeout为1,单位为天
joinExecutor.awaitTermination(1, TimeUnit.DAYS);
}
//自定义方法exeAsync,传入一个javabean
public void exeAsync(MyBean myBean){
CompletableFuture.runAsync(() -> {
try{
//自定义方法myFunc,传入参数myBean
myFunc(myBean);
} catch(Exception e){
log.error(e.getMessage(), e);
}
}, joinExecutor);
}
//还有一种方法
public void exeAsync2(MyBean myBean){
joinExecutor.execute(() -> {
try{
//自定义方法myFunc,传入参数myBean
myFunc(myBean);
} catch(Exception e){
log.error(e.getMessage(), e);
}
});
}
}
三、kafka常见问题与解决方法
1.kafka不用在控制台手动创建队列,直接使用即可;但是使用新队列时,需要生产者先向该队列发送一条消息,然后消费者才能启动成功;如果消费者先启动并监听新队列,会报错(队列不存在导致):
nested exception is java.lang.IllegalStateException: Topic(s) [test-my_kafka_Topic] is/are not present and missingTopicFatal is true
如果必须先启动消费者,那要在yml的kafka配置中增加:
missing-topics-fatal: false
详细yml例子如下:
spring:
kafka:
consumer:
group-id: xxxGroup
bootstrap-servers: 10.123.123.123:9092
listener:
missing-topics-fatal: false
更多推荐
所有评论(0)