智能BI项目-接口的异步化(六)
往期开发(已经完成)
- 智能 BI 项目-介绍(一)
- 智能 BI 项目-初始化(二)
- 智能 BI 项目-初学 AI 分析(三)
- 智能 BI 项目-AI 接口调用(四)
- 智能 BI 项目-接口优化(五)
- 智能 BI 项目-接口的异步化(六)
- 智能 BI 项目-引入 RabbitMQ(七)
今日后端开发过程
系统问题分析
问题
用户等待的智能分析时间过长,由于要调用第三方 AI 服务 等待时间过长
业务服务器可能会有很多请求在处理,导致服务器资源紧张,压力过大时候会导致服务器宕机或者无法处理新的请求。
调用第三方 AI 服务,AI 处理能力是有限的,例如,每 3s 只能处理一个请求,会导致 AI 处理压力多大,处理不过来的情况,严重时候 AI 可能会对服务器系统拒绝服务
解决方案
异步化接口,原来的同步化与异步化的区别
同步:一件事做完,才能继续做下件事情。比如加工需要等上一个零件弄完才能继续下一个
异步:不需要等待上一件事情做完,可以先做记录或者是保存当前任务,等待任务队列进行分配,此时可以做其他事情,等待这个事件做完,通知用户,可以继续做后续事情。
业务流程分析
标准化分析
当用户要进行等待时间或是耗时很长的操作室,点击提交之后,不需要在页面进行等待,而是将此次任务提交到后端保存至数据库中。
用户点击提交新任务时:
a. 任务提交成功:
i . 如果程序中还有多余的空闲线程,可以立即去处理这个任务。
ii. 如果程序中的线程都在处理任务中,无法处理此任务,放到等待队列中。
b. 任务提交失败(所有的线程都在处理并且任务队列满了):
i . 拒绝掉这个任务,再也不去执行。
ii. 将此次任务保存到数据库中,来记录本次任务失败的情况,并在程序的线程空闲时候,可以把任务从数据库中取出再去执行。程序的线程从任务队列中取出任务依次执行,每完成一件事情要修改一下任务的状态。
用户可以查询当前任务的执行状态,在任务执行成功或者失败时候可以收到通知(发送 email,系统的消息提示,手机短信等),从而优化用户体验。
如果我们要执行的任务非常的复杂,包含很多个环节,在每个小任务完成时,程序需要记录一下任务执行的状态(进度)。
智能 BI 业务分析
用户点击提交按钮时候,先将要分析的信息保存到数据库中作为一个任务。
用户可以我的图表页面查看所有的图表(状态:已生成的。生成中、生成失败)的信息和状态。
用户可以修改生成失败的图表信息点击重新生成。
优化前的请求流程图:

异步优化后的请求流程图:

优化后出现新的问题:
a.任务队列的最大容量应该设置为多少?
b.程序如何从任务队列中取出任务去执行?
c.这个任务队列的流程如何在程序中实现?
d.如何保证程序中最多同时执行多少个任务?解决方案-线程池
为什么要引入线程池?
- 线程的管理比较复杂(例如,何时增加线程,何时减少空闲线程)
- 任务存取比较复杂(何时接受任务,何时拒绝任务,如何保证执行时不抢到同一个任务)
线程池的作用
- 轻松帮助管理线程池
- 协调任务的执行和分配过程
简单流程图如下

线程池的实现
自行不用实现,如果是在 Spring 框架中可以使用 ThreadPoolTaskExecutor 配合注解@Async 来实现(不建议)。
如果是在 Java 中,可以使用 JUC 并发编程包中的 ThreadPoolExecutor 来实现非常灵活地自定义线程池,通用性更加高。线程池的参数如何设置?
1
2
3
4
5
6
7public ThreadPoolExecutor(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue,
ThreadFactory threadFactory,
RejectedExecutionHandler handler){...}线程池参数理解
- corePoolSize(核心线程数 -> 可以当做是公司的正式员工,也就是常驻员工):正常情况下,系统应该能保持同时工作的线程数(随时属于就绪状态)。
- maximumPoolSize(最大线程数 -> 可以当做是公司人手不够,最多能招的人数):极端情况下,线程池最多能开启多少线程(包括正式员工和临时员工)。
- keepAliveTime(空闲线程存活时间 -> 可以当做是非正式员工的工作时间):非核心线程在没有任务的情况下,过多久要删除(工作时长),释放无用线程资源。
- TimeUnit(空闲线程的存活时间单位):min、s
- workQueue(任务队列):用户存放给线程执行的任务队列,存在一个任务队列长度(一定要进行设置,不能让任务队列长度无限,不然则占用资源过大)
- ThreadFactory(线程工厂):控制每个线程的生成,以及线程的一些参数设置(比如定义线程名字)
- RejectedExecutionHandler(拒绝策略):当线程数最大且都在处理任务时任务队列此时也已经占满,如果再来任务就要采取措施,比如抛出异常,自定义处理策略等。
如何确定线程池的参数?需要结合具体的业务场景,不断的优化和设置。
- 假设 AI 生成能力的并发只允许 4 个任务同时执行,AI 能力允许 20 个任务去排队
- corePoolSize(核心线程数 -> 正式员工数):正常情况下可是设置为 2 - 4
- maximumPoolSize(极限线程数 -> 正式员工数 + 临时员工数) 设置为 <= 4
- keepAliveTime(空闲线程存活时间): 一般设置为 min 级、s 级
- TimeUnit(空闲线程存活时间单位): min、s
- workQueue(工作队列):结合实际情况设置,可以设置为 20
- threadFactory(线程工厂):线程的生成和属性的设置
- rejectedExecutionHandler(拒绝策略):抛出异常,标记数据库的任务状态为”任务满了已拒绝”
一般情况下,任务分为 IO 密集型和计算密集型(CPU 密集型)两种:
计算密集型:CPU 性能,比如音频处理,图像处理,数学计算等,一般设置 corePoolSize 为 CPU 的核数 + 1,多出一恶搞是为了让每个线程都能利用好 CPU 的每个核,而且避免线程频繁切换(减少争抢,减少开销)
IO 密集型:吃带宽/内存/硬盘的读写资源,corePoolSize 可以设置大一些,一般是 2N 左右,但是建议以 IO 的能力为主。线城池的工作机制(流程图)
- 刚开始没有任何线程和任务

- 此时来了一个任务,公司发现正式员工还没有占满(corePoolSize = 2),喊来一个正式员工(核心线程)来处理任务

- 此时又来了一个任务,公司发现正式员工还没有占满(corePoolSize = 2),再喊来一个正式员工(核心线程)来处理任务

- 此时又来了两个任务,公司发现正式员工还已经占满(corePoolSize = 2),但是并不会去招临时工(添加新的线程)来处理,而是先放入任务队列(workQueue.size=2)中。

- 此时又来了两个任务,公司发现正式员工还已经占满(corePoolSize = 2),且任务队列已经占满了(当前线程数 > workQueue.size = 2,队列中已有任务出=当前最大任务数),新增临时员工(最大线程 maximumPoolSize=4,未达到极限线程数)来处理新的任务,而不是丢弃任务

- 此时又来了一个任务,公司发现正式员工还已经占满(corePoolSize = 2),且任务队列已经占满了(当前线程数 > workQueue.size = 2,队列中已有任务出=当前最大任务数),并且此时 maximumPoolSize=4,已经达到了极限线程数,此时的任务要么被拒绝、要么自定义策略进行处理(RejectedExecutionHandler)

- 如果当线程数超过 corePoolSize(正式员工数),又没有新任务给他,那么等到 keepAliveTime 时间达到后会释放非核心线程(临时员工)
- 刚开始没有任何线程和任务
代码开发
- 自定义线程池
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
public class ThreadPoolExecutorConfig {
public ThreadPoolExecutor threadPoolExecutor() {
ThreadFactory threadFactory = new ThreadFactory() {
private int count = 1;
public Thread newThread( Runnable r) {
Thread thread = new Thread(r);
thread.setName("线程" + count);
count++;
return thread;
}
};
RejectedExecutionHandler rejectedExecutionHandler = new RejectedExecutionHandler() {
public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) {
throw new BusinessException(ErrorCode.TOO_MANY_REQUEST, "网络繁忙,请稍后再提交任务");
}
};
return new ThreadPoolExecutor(2, 4, 100, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(4), threadFactory, rejectedExecutionHandler);
}
}- 定义 controller 进行测试
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
private ThreadPoolExecutor threadPoolExecutor;
public void add( String name) {
CompletableFuture.runAsync(() -> {
System.out.println(name + "-任务执行中,当前线程为:" + Thread.currentThread().getName());
try {
Thread.sleep(600000);
} catch (InterruptedException e) {
e.printStackTrace();
}
}, threadPoolExecutor);
}
public String get() {
Map<String, Object> map = new HashMap<>();
int size = threadPoolExecutor.getQueue().size();
map.put("队列长度", size);
long taskCount = threadPoolExecutor.getTaskCount();
map.put("任务总数", taskCount);
long completedTaskCount = threadPoolExecutor.getCompletedTaskCount();
map.put("已完成任务数", completedTaskCount);
int activeCount = threadPoolExecutor.getActiveCount();
map.put("正在工作的线程数", activeCount);
return JSONUtil.toJsonStr(map);
}完整代码实现
- 给表中添加新的字段,status-状态,以及 execMessage 消息
- 用户点击提交任务时,先把图表信息立即保存到数据库中,然后提交任务
- 用户提交任务之后,现将状态修改为“执行中”,执行成功之后修改为“已完成”,保存执行结果;执行失败后,状态修改为失败,记录任务失败信息
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89/**
* AI异步对话分析
*
* @param multipartFile
* @param genChartByAiRequest
* @param request
* @return
*/
public BaseResponse<String> genChartByAIUseAsync( MultipartFile multipartFile,
GenChartByAiRequest genChartByAiRequest, HttpServletRequest request) {
// 参数获取
String goal = genChartByAiRequest.getGoal();
String name = genChartByAiRequest.getName();
String chartType = genChartByAiRequest.getChartType();
// 参数校验
ThrowUtils.throwIf(StringUtils.isEmpty(goal), ErrorCode.PARAMS_ERROR, "目标不能为空");
ThrowUtils.throwIf(StringUtils.isEmpty(chartType), ErrorCode.PARAMS_ERROR, "类型不能为空");
ThrowUtils.throwIf(StringUtils.isEmpty(name), ErrorCode.PARAMS_ERROR, "图表名称不能为空");
// 文件校验
long size = multipartFile.getSize();
String originalFilename = multipartFile.getOriginalFilename();
ThrowUtils.throwIf(size > MAX_FILE_SIZE, ErrorCode.PARAMS_ERROR, "文件过大不得超过1MB");
ThrowUtils.throwIf(!FILE_NAME_LIST.contains(FileUtil.getSuffix(originalFilename)), ErrorCode.PARAMS_ERROR, "不支持此文件");
// 判断是否登录
User loginUser = userService.getLoginUser(request);
ThrowUtils.throwIf(Objects.isNull(loginUser), ErrorCode.NOT_LOGIN_ERROR);
redisLimiterManager.doRateLimit("genChartByAI_" + loginUser.getId());
// 用户消息拼接
StringBuilder userMessageBuilder = new StringBuilder();
userMessageBuilder.append("分析需求:").append("\n");
String newGoal = goal + ",请使用" + chartType;
userMessageBuilder.append(newGoal).append("\n");
userMessageBuilder.append("原始数据:").append("\n");
String dataStr = ExcelUtils.excelToCsv(multipartFile);
userMessageBuilder.append(dataStr);
// 先保存任务到数据库中
Chart chart = new Chart();
chart.setName(name);
chart.setGoal(goal);
chart.setChartData(dataStr);
chart.setChartType(chartType);
chart.setStatus("wait");
chart.setExecMessage("任务等待中");
chart.setUserId(loginUser.getId());
boolean firstChartSave = chartService.save(chart);
ThrowUtils.throwIf(!firstChartSave, ErrorCode.PARAMS_ERROR, "任务提交失败");
// 将任务转入任务队列中分配执行
CompletableFuture.runAsync(() -> {
Chart updateFirstChart = new Chart();
updateFirstChart.setId(chart.getId());
updateFirstChart.setStatus("running");
updateFirstChart.setExecMessage("任务正在执行中");
boolean isFirstUpdate = chartService.updateById(updateFirstChart);
if (!isFirstUpdate) {
handleChartUpdateError(chart.getId());
return;
}
// AI接口服务
String content = aiManager.doChat(userMessageBuilder.toString());
String[] splits = content.split("【【【【【");
if (splits.length != 3) {
throw new RuntimeException("AI 生成错误");
}
Chart updateChartData = new Chart();
updateChartData.setId(chart.getId());
updateChartData.setStatus("succeed");
updateFirstChart.setExecMessage("任务已经完成");
updateChartData.setGenResult(splits[2].trim());
updateChartData.setGenChart(splits[1].trim());
boolean isUpdateChartData = chartService.updateById(updateChartData);
if (!isUpdateChartData) {
handleChartUpdateError(chart.getId());
return;
}
}, threadPoolExecutor);
return ResultUtils.success("任务提交成功!");
}
private void handleChartUpdateError(long chartId) {
Chart updateChartResult = new Chart();
updateChartResult.setId(chartId);
updateChartResult.setStatus("failed");
updateChartResult.setExecMessage("更新图表状态失败!");
boolean updateResult = chartService.updateById(updateChartResult);
if (!updateResult) {
log.error("更新图表失败状态失败" + chartId + "," + "更新图表状态失败!");
}
}
后端收获
今天主要学习了 JUC 并发编程,使用线程池来解决接口响应时间较长,等待时间长的问题,由于以前的项目中很少使用线程,此次也是对线程池有了一个清晰的理解,以及线程池的简单使用,并写入到接口中,实际测试,相对于来说收获还是蛮多的,接触新的知识花的时间也多,为了就是更好的、更深入的学习并发编程的奥妙。
今日前端的开发
- 添加新的路由信息
1
{ path: '/chart-add/async', name: '制作图表(异步)', icon: 'lineChart', component: './Chart/AsyncAddChart' },
- 添加新的路由信息
- 优化原来的页面,接入新的接口,重新调用 openAPI 更新接口代码
- 前端的优化代码如下:
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94const asyncAddChart: React.FC = () => {
const [submitting, setSubmitting] = useState<boolean>(false);
const [form] = useForm();
const normFile = (e: any) => {
console.log('Upload event:', e);
if (Array.isArray(e)) {
return e;
}
return e?.fileList;
};
const onFinish = async (values: any) => {
if (submitting) return;
setSubmitting(true);
const params = {
...values,
file: undefined,
};
try {
const res = await genChartByAIUseAsyncUsingPOST(params, {}, values?.file[0]?.originFileObj);
if (res?.data) {
message.success(res.data);
form.resetFields();
}
} catch (e: any) {
message.error('分析失败!');
}
setSubmitting(false);
};
return (
<div className="async-add-chart">
<Card title={'智能分析(异步)'}>
<Form form={form} name="addChart" onFinish={onFinish} labelAlign={'left'}>
<Form.Item
name="name"
label="图表名称:"
rules={[{ required: true, message: '请输入图表的名称' }]}
>
<Input placeholder="请填写图表名称" />
</Form.Item>
<Form.Item
name="goal"
label="分析诉求:"
rules={[{ required: true, message: '请输入诉求' }]}
>
<TextArea placeholder="请填写您的诉求" />
</Form.Item>
<Form.Item
name="chartType"
label="图表类型:"
hasFeedback
rules={[{ required: true, message: '请选择图表类型' }]}
>
<Select
placeholder="请选择图表类型"
options={[
{ value: '折线图', label: '折线图' },
{ value: '柱状图', label: '柱状图' },
{ value: '堆叠图', label: '堆叠图' },
{ value: '饼图', label: '饼图' },
{ value: '雷达图', label: '雷达图' },
]}
></Select>
</Form.Item>
<Form.Item
rules={[{ required: true, message: '请上传文件' }]}
name="file"
label="上传文件:"
valuePropName="fileList"
getValueFromEvent={normFile}
>
<Upload name="file">
<Button icon={<UploadOutlined />}>数据文件</Button>
</Upload>
</Form.Item>
<Form.Item wrapperCol={{ span: 12, offset: 3 }}>
<Space>
<Button type="primary" htmlType="submit" loading={submitting} disabled={submitting}>
{submitting ? '分析中……' : '提交分析(异步)'}
</Button>
<Button htmlType="reset" disabled={submitting}>
重置
</Button>
</Space>
</Form.Item>
</Form>
</Card>
</div>
);
};
export default asyncAddChart;效果图如下
- 填写图表

- 提交查看

- 结果渲染

- 填写图表
- 前端收获
如何渲染异步化的组件显示,熟悉 ant design 组件库,“Result” 组件的使用。

