增加drui连接池配置&修改失败补偿批次查询&脱敏规则增加redis缓存

This commit is contained in:
zren25
2025-05-27 14:30:24 +08:00
parent eb860707b6
commit 390047266d
14 changed files with 361 additions and 46 deletions

View File

@@ -61,6 +61,7 @@
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>druid-spring-boot-starter</artifactId>
<version>1.2.18</version>
</dependency>
<dependency>
<groupId>io.springfox</groupId>
@@ -71,7 +72,7 @@
<dependency>
<groupId>com.baomidou</groupId>
<artifactId>mybatis-plus-boot-starter</artifactId>
<version>3.3.0</version>
<version>3.5.3.1</version>
</dependency>
<dependency>
@@ -157,11 +158,6 @@
<groupId>com.zaxxer</groupId>
<artifactId>HikariCP</artifactId>
</dependency>
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>druid-spring-boot-starter</artifactId>
<version>1.1.13</version>
</dependency>
<!--volvo-->
<dependency>
<groupId>com.volvo</groupId>

View File

@@ -1,21 +1,21 @@
package com.volvo.ai.analytic.center;
import com.volvo.common.feign.annotation.EnableVolvoFeignClients;
import org.mybatis.spring.annotation.MapperScan;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration;
import org.springframework.cloud.client.discovery.EnableDiscoveryClient;
import org.springframework.context.annotation.Bean;
import org.springframework.retry.annotation.EnableRetry;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.web.client.RestTemplate;
import com.volvo.common.feign.annotation.EnableVolvoFeignClients;
@SpringBootApplication(scanBasePackages = {"com.volvo"})
@SpringBootApplication(scanBasePackages = {"com.volvo"},exclude = {DataSourceAutoConfiguration.class})
@EnableVolvoFeignClients
@EnableScheduling
@MapperScan({"com.volvo.ai.analytic.center.mapper"})
@MapperScan("com.volvo.ai.analytic.center.mapper")
@EnableDiscoveryClient
@EnableRetry
public class AiAnalyticCenterServiceApplication {

View File

@@ -22,7 +22,7 @@ public class ExecutorConfig {
@Value("${task.pool.queueCapacity}")
private int queueCapacity;
@Bean("corpusProcessExecutor")
@Bean("threadPoolTaskExecutor")
public ThreadPoolTaskExecutor corpusProcessExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
int corePoolSize = Runtime.getRuntime().availableProcessors() + 2;

View File

@@ -0,0 +1,24 @@
package com.volvo.ai.analytic.center.config;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.data.redis.serializer.GenericJackson2JsonRedisSerializer;
import org.springframework.data.redis.serializer.StringRedisSerializer;
@Configuration
public class RedisConfig {
@Bean
public RedisTemplate<String, Object> redisTemplate(RedisConnectionFactory factory) {
RedisTemplate<String, Object> template = new RedisTemplate<>();
template.setConnectionFactory(factory);
template.setKeySerializer(new StringRedisSerializer());
template.setValueSerializer(new GenericJackson2JsonRedisSerializer());
template.setHashKeySerializer(new StringRedisSerializer());
template.setHashValueSerializer(new GenericJackson2JsonRedisSerializer());
template.afterPropertiesSet();
return template;
}
}

View File

@@ -14,10 +14,9 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.kafka.core.KafkaTemplate;
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.bind.annotation.RestController;
import org.springframework.web.bind.annotation.*;
import javax.sql.DataSource;
@RestController
@@ -41,6 +40,9 @@ public class TestController {
@Qualifier("dccKafkaTemplate")
private KafkaTemplate<String, String> dccKafkaProducer;
@Autowired
private DataSource dataSource;
@PostMapping("/mockMq")
@ApiOperation(value = "补偿处理消息")
public ResultMsg<Object> mockMq(@RequestBody String message) {
@@ -75,5 +77,10 @@ public class TestController {
return ResultMsg.ok("ok");
}
@GetMapping("/pool")
public String checkPool() {
return dataSource.getClass().getName();
}
}

View File

@@ -32,6 +32,8 @@ import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource;
import java.util.Arrays;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
@Slf4j
@@ -89,13 +91,17 @@ public class CorpusFailJob {
List<AiAnalysisErrors> aiAnalysisErrorsListlist = aiAnalysisErrorsService.queryAnalysisErrorList(Arrays.asList(BusinessTypeEnum.SMART_ASSISTANT.getCode(),BusinessTypeEnum.SMART_ASSISTANT_QIWEI.getCode(),BusinessTypeEnum.SMART_ASSISTANT_NAMEPLATE.getCode()),offset, pageSize);
if(CollectionUtils.isNotEmpty(aiAnalysisErrorsListlist)) {
log.info("corpusFailTask语料解析失败重试处理 size:{}", aiAnalysisErrorsListlist.size());
aiAnalysisErrorsListlist.stream().forEach(aiAnalysisErrors -> {
List<String> requestIdList = aiAnalysisErrorsListlist.stream()
.map(AiAnalysisErrors::getAiAnalysisRequestId)
.collect(Collectors.toList());
List<AiAnalysisRequestLogs> aiAnalysisRequestLogsList = aiAnalysisRequestLogsMapper.queryByAiAnalysisRequestIds(requestIdList);
Map<String, AiAnalysisRequestLogs> requestLogsMap = aiAnalysisRequestLogsList.stream()
.collect(Collectors.toMap( AiAnalysisRequestLogs::getAiAnalysisRequestId,logs -> logs ));
aiAnalysisErrorsListlist.forEach(aiAnalysisErrors -> {
try {
LambdaQueryWrapper<AiAnalysisRequestLogs> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisErrors.getAiAnalysisRequestId());
AiAnalysisRequestLogs oldAiAnalysisRequestLogs = aiAnalysisRequestLogsMapper.selectOne(queryWrapper);
AiAnalysisRequestLogs oldAiAnalysisRequestLogs = requestLogsMap.get(aiAnalysisErrors.getAiAnalysisRequestId());
if (null != oldAiAnalysisRequestLogs && StringUtils.isBlank(oldAiAnalysisRequestLogs.getBusinessResponse())) {
DiFyReq diFyReq = JSONObject.parseObject(oldAiAnalysisRequestLogs.getDifyRequest(), DiFyReq.class);
@@ -203,11 +209,16 @@ public class CorpusFailJob {
List<AiAnalysisErrors> aiAnalysisErrorsListlist = aiAnalysisErrorsService.queryAnalysisErrorList(Arrays.asList(BusinessTypeEnum.INTELLIGENT_CUSTOMER.getCode()),offset, pageSize);
if(CollectionUtils.isNotEmpty(aiAnalysisErrorsListlist)) {
log.info("intelligentCustomer4in1FailTask语料解析失败重试处理 size:{}", aiAnalysisErrorsListlist.size());
aiAnalysisErrorsListlist.stream().forEach(aiAnalysisErrors -> {
List<String> requestIdList = aiAnalysisErrorsListlist.stream()
.map(AiAnalysisErrors::getAiAnalysisRequestId)
.collect(Collectors.toList());
List<AiAnalysisRequestLogs> aiAnalysisRequestLogsList = aiAnalysisRequestLogsMapper.queryByAiAnalysisRequestIds(requestIdList);
Map<String, AiAnalysisRequestLogs> requestLogsMap = aiAnalysisRequestLogsList.stream()
.collect(Collectors.toMap( AiAnalysisRequestLogs::getAiAnalysisRequestId,logs -> logs ));
aiAnalysisErrorsListlist.forEach(aiAnalysisErrors -> {
try {
LambdaQueryWrapper<AiAnalysisRequestLogs> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisErrors.getAiAnalysisRequestId());
AiAnalysisRequestLogs oldAiAnalysisRequestLogs = aiAnalysisRequestLogsMapper.selectOne(queryWrapper);
AiAnalysisRequestLogs oldAiAnalysisRequestLogs = requestLogsMap.get(aiAnalysisErrors.getAiAnalysisRequestId());
if (null != oldAiAnalysisRequestLogs) {
JSONObject businessRequest = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest());

View File

@@ -1,5 +1,5 @@
package com.volvo.ai.analytic.center.mapper;
import java.util.List;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
import org.apache.ibatis.annotations.Mapper;
@@ -9,4 +9,8 @@ import org.apache.ibatis.annotations.Param;
public interface AiAnalysisRequestLogsMapper extends BaseMapper<AiAnalysisRequestLogs> {
public AiAnalysisRequestLogs queryAiAnalysisRequestLogsByBusinessReponse(@Param("sourceId") String sourceId);
List<AiAnalysisRequestLogs> queryByAiAnalysisRequestIds(@Param("aiAnalysisRequestIds") List<String> aiAnalysisRequestIds);
AiAnalysisRequestLogs queryByAiAnalysisRequestId(@Param("aiAnalysisRequestId") String aiAnalysisRequestId);
}

View File

@@ -2,8 +2,8 @@ package com.volvo.ai.analytic.center.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.volvo.ai.analytic.center.entity.DataMaskingRule;
import org.springframework.stereotype.Repository;
import org.apache.ibatis.annotations.Mapper;
@Repository
@Mapper
public interface DataMaskingRuleMapper extends BaseMapper<DataMaskingRule> {
}

View File

@@ -2,8 +2,8 @@ package com.volvo.ai.analytic.center.mapper;
import com.baomidou.mybatisplus.core.mapper.BaseMapper;
import com.volvo.ai.analytic.center.entity.DiffdefeatApprove;
import org.springframework.stereotype.Repository;
import org.apache.ibatis.annotations.Mapper;
@Repository
@Mapper
public interface DiffdefeatApproveMapper extends BaseMapper<DiffdefeatApprove> {
}

View File

@@ -20,6 +20,7 @@ import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.RestController;
@@ -27,7 +28,9 @@ import javax.annotation.Resource;
import java.time.LocalDateTime;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.*;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
/**
* @ClassName CorpusProcessKafkaConsumer
@@ -54,8 +57,9 @@ public class CorpusProcessKafkaProducer {
@Resource
private RocketMQTemplate rocketMqTemplate;
@Resource(name = "corpusProcessExecutor")
private ExecutorService executor;
@Autowired
@Resource(name = "threadPoolTaskExecutor")
private ThreadPoolTaskExecutor executor;
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}")
public void listen(List<ConsumerRecord<String, Object>> recordMessages) {
@@ -82,9 +86,10 @@ public class CorpusProcessKafkaProducer {
// 2. 安全关闭线程池
executor.shutdown();
try {
if (!executor.awaitTermination(60, TimeUnit.SECONDS)) {
ThreadPoolExecutor threadPoolExecutor = this.executor.getThreadPoolExecutor();
if (!threadPoolExecutor.awaitTermination(60, TimeUnit.SECONDS)) {
log.warn("Thread pool did not terminate in time. Forcing shutdown.");
executor.shutdownNow();
threadPoolExecutor.shutdownNow();
}
} catch (InterruptedException e) {
log.error("Thread pool termination interrupted: ", e);

View File

@@ -20,13 +20,14 @@ public class AiAnalysisRequestLogsServiceImpl extends ServiceImpl<AiAnalysisRequ
@Override
public boolean saveAiAnalysisRequestLogs(AiAnalysisRequestLogs aiAnalysisRequestLogs) {
LambdaQueryWrapper<AiAnalysisRequestLogs> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisRequestLogs.getAiAnalysisRequestId());
AiAnalysisRequestLogs oldAiAnalysisRequestLogs= aiAnalysisRequestLogsMapper.selectOne(queryWrapper);
AiAnalysisRequestLogs oldAiAnalysisRequestLogs= aiAnalysisRequestLogsMapper.queryByAiAnalysisRequestId(aiAnalysisRequestLogs.getAiAnalysisRequestId());
if (oldAiAnalysisRequestLogs == null) {
aiAnalysisRequestLogs.setCreateTime(new Date());
return aiAnalysisRequestLogsMapper.insert(aiAnalysisRequestLogs) > 0;
} else {
LambdaQueryWrapper<AiAnalysisRequestLogs> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
aiAnalysisRequestLogs.setUpdateTime(new Date());
return aiAnalysisRequestLogsMapper.update(aiAnalysisRequestLogs, queryWrapper) > 0;
}
@@ -34,9 +35,7 @@ public class AiAnalysisRequestLogsServiceImpl extends ServiceImpl<AiAnalysisRequ
@Override
public AiAnalysisRequestLogs queryByAiAnalysisRequestId(String aiAnalysisRequestId) {
LambdaQueryWrapper<AiAnalysisRequestLogs> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisRequestId);
return aiAnalysisRequestLogsMapper.selectOne(queryWrapper);
return aiAnalysisRequestLogsMapper.queryByAiAnalysisRequestId(aiAnalysisRequestId);
}
@Override

View File

@@ -1,5 +1,6 @@
package com.volvo.ai.analytic.center.service.impl;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.baomidou.mybatisplus.extension.service.impl.ServiceImpl;
import com.volvo.ai.analytic.center.constant.YesOrNoConstants;
import com.volvo.ai.analytic.center.dto.req.RunMaskingRuleInput;
@@ -7,27 +8,56 @@ import com.volvo.ai.analytic.center.entity.DataMaskingRule;
import com.volvo.ai.analytic.center.enums.RuleCategoryEnum;
import com.volvo.ai.analytic.center.mapper.DataMaskingRuleMapper;
import com.volvo.ai.analytic.center.service.DataMaskingRuleService;
import com.volvo.ai.analytic.center.utils.RedisUtil;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Service;
import java.util.Collections;
import java.util.List;
import java.util.Objects;
import java.util.concurrent.TimeUnit;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
@Slf4j
@Service
public class DataMaskingRuleServiceImpl extends ServiceImpl<DataMaskingRuleMapper, DataMaskingRule> implements DataMaskingRuleService {
@Autowired
private RedisUtil redisUtil;
private static final String CACHE_KEY_PREFIX = "analtyticCenter:data_masking_rule:";
@Autowired
private DataMaskingRuleMapper dataMaskingRuleMapper;
@Override
public List<DataMaskingRule> getDataMaskingRuleListByApplicationChannel(String applicationChannel) {
log.info("getDataMaskingRuleListByApplicationChannel {}", applicationChannel);
//根据适用渠道applicationChannel获取数据脱敏规则List
List<DataMaskingRule> dataMaskingRuleList = this.lambdaQuery()
.like(DataMaskingRule::getApplicationChannel, applicationChannel)
.eq(DataMaskingRule::getRuleStatus, YesOrNoConstants.YES)
.eq(DataMaskingRule::getIsDeleted, YesOrNoConstants.NO)
.list();
return dataMaskingRuleList;
String cacheKey = CACHE_KEY_PREFIX + applicationChannel;
// 1. 先从缓存获取
List<DataMaskingRule> cachedRules = redisUtil.lrange(cacheKey, 0, -1, DataMaskingRule.class);
if (cachedRules == null || cachedRules.isEmpty()) {
log.info("getDataMaskingRuleListByApplicationChannel 查询数据库: {}", applicationChannel);
// 2. 查询数据库
cachedRules = dataMaskingRuleMapper.selectList(
new LambdaQueryWrapper<DataMaskingRule>()
.like(DataMaskingRule::getApplicationChannel, applicationChannel)
.eq(DataMaskingRule::getRuleStatus, YesOrNoConstants.YES)
.eq(DataMaskingRule::getIsDeleted, YesOrNoConstants.NO)
);
// 3. 写入缓存
if (cachedRules != null && !cachedRules.isEmpty()) {
redisUtil.rpush(cacheKey, cachedRules.toArray());
redisUtil.expire(cacheKey, 1, TimeUnit.DAYS); // 设置过期时间
}
} else {
log.info("getDataMaskingRuleListByApplicationChannel 从缓存获取: {}", applicationChannel);
}
return cachedRules;
}
@Override

View File

@@ -0,0 +1,204 @@
package com.volvo.ai.analytic.center.utils;
import java.util.Collections;
import java.util.List;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Component;
import javax.annotation.Resource;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors;
/**
* Redis 工具类,封装常用的 Redis 操作
*/
@Component
public class RedisUtil {
@Resource
private RedisTemplate<String, Object> redisTemplate;
/**
* 设置缓存
*
* @param key 键
* @param value 值
*/
public void set(String key, Object value) {
redisTemplate.opsForValue().set(key, value);
}
/**
* 设置缓存并指定过期时间
*
* @param key 键
* @param value 值
* @param timeout 时间
* @param unit 时间单位
*/
public void set(String key, Object value, long timeout, TimeUnit unit) {
redisTemplate.opsForValue().set(key, value, timeout, unit);
}
/**
* 获取缓存
*
* @param key 键
* @return 值
*/
public Object get(String key) {
return key == null ? null : redisTemplate.opsForValue().get(key);
}
/**
* 删除缓存
*
* @param key 键
*/
public void delete(String key) {
if (key != null) {
redisTemplate.delete(key);
}
}
/**
* 判断缓存是否存在
*
* @param key 键
* @return 是否存在
*/
public boolean hasKey(String key) {
return Boolean.TRUE.equals(redisTemplate.hasKey(key));
}
/**
* 设置过期时间
*
* @param key 键
* @param timeout 过期时间
* @param unit 时间单位
* @return 是否成功
*/
public boolean expire(String key, long timeout, TimeUnit unit) {
return Boolean.TRUE.equals(redisTemplate.expire(key, timeout, unit));
}
/**
* 获取剩余过期时间
*
* @param key 键
* @param unit 返回的时间单位
* @return 剩余时间
*/
public Long getExpire(String key, TimeUnit unit) {
return redisTemplate.getExpire(key, unit);
}
// ==================== Map ======================
/**
* HashGet
*/
public Object hget(String key, String item) {
return redisTemplate.opsForHash().get(key, item);
}
/**
* HashSet
*/
public void hset(String key, String item, Object value) {
redisTemplate.opsForHash().put(key, item, value);
}
// ==================== List ======================
/**
* 获取 list 缓存的内容
*
* @param key 键
* @param start 开始
* @param end 结束
* @return 列表范围 [start, end]
*/
public List<Object> lrange(String key, long start, long end) {
return redisTemplate.opsForList().range(key, start, end);
}
public <T> List<T> lrange(String key, long start, long end, Class<T> clazz) {
List<Object> rawList = redisTemplate.opsForList().range(key, start, end);
if (rawList == null) return Collections.emptyList();
return rawList.stream()
.filter(clazz::isInstance)
.map(clazz::cast)
.collect(Collectors.toList());
}
/**
* 获取 list 缓存的长度
*/
public Long llen(String key) {
return redisTemplate.opsForList().size(key);
}
/**
* 将多个值添加到列表头部
*/
public Long lpush(String key, Object... values) {
return redisTemplate.opsForList().leftPushAll(key, values);
}
/**
* 将多个值添加到列表尾部
*/
public Long rpush(String key, Object... values) {
return redisTemplate.opsForList().rightPushAll(key, values);
}
/**
* 移除并返回列表的第一个元素
*/
public Object lpop(String key) {
return redisTemplate.opsForList().leftPop(key);
}
/**
* 移除并返回列表的最后一个元素
*/
public Object rpop(String key) {
return redisTemplate.opsForList().rightPop(key);
}
// ==================== Set ======================
/**
* 向集合中添加成员
*/
public Long sadd(String key, Object... values) {
return redisTemplate.opsForSet().add(key, values);
}
/**
* 获取集合中的所有成员
*/
public List<Object> smembers(String key) {
return (List<Object>) redisTemplate.opsForSet().members(key);
}
// ==================== ZSet ======================
/**
* 添加一个有序集合成员
*/
public Boolean zadd(String key, Object value, double score) {
return redisTemplate.opsForZSet().add(key, value, score);
}
/**
* 获取有序集合中的成员
*/
public Set<Object> zrange(String key, long start, long end) {
return redisTemplate.opsForZSet().range(key, start, end);
}
}

View File

@@ -16,4 +16,39 @@
AND JSON_EXTRACT( CAST( business_response AS JSON ), '$.recordId' ) = #{sourceId} limit 1
</select>
<select id="queryByAiAnalysisRequestIds" resultType="com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs" >
SELECT
id,
`ai_analysis_request_id` AS aiAnalysisRequestId,
`ai_analysis_request_type` AS aiAnalysisRequestType,
`business_request` AS businessRequest,
`dify_request` AS difyRequest,
`dify_response` AS difyResponse,
`business_response` AS businessTesponse
FROM
tt_ai_analysis_request_logs
WHERE is_deleted='0'
and ai_analysis_request_id in
<foreach collection="aiAnalysisRequestIds" item="aiAnalysisRequestId" open="(" separator="," close=")">
#{aiAnalysisRequestId}
</foreach>
</select>
<select id="queryByAiAnalysisRequestId" resultType="com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs" >
SELECT
id,
`ai_analysis_request_id` AS aiAnalysisRequestId,
`ai_analysis_request_type` AS aiAnalysisRequestType,
`business_request` AS businessRequest,
`dify_request` AS difyRequest,
`dify_response` AS difyResponse,
`business_response` AS businessTesponse
FROM
tt_ai_analysis_request_logs
WHERE is_deleted='0'
and ai_analysis_request_id = #{aiAnalysisRequestId} limit 1
</select>
</mapper>