Merge branch 'feature_20250519_nameplate' into dev-feature-20250903-Portrait0603

# Conflicts:
#	ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/TestController.java
#	ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/CorpusFailJob.java
#	ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/job/NameplateCorpusJob.java
#	ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/mq/NameplateKafkaConsumer.java
#	ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/service/impl/TmNameplateCorpusServiceImpl.java
This commit is contained in:
ZLI263
2025-09-04 14:47:06 +08:00
12 changed files with 80 additions and 176 deletions

View File

@@ -23,7 +23,9 @@ public class KafkaConfig {
// 第一个Kafka配置
/**
* 湖仓Kafka配置用途mock数据
*/
@Bean(name = "dccKafkaTemplate")
public KafkaTemplate<String, String> dccKafkaTemplate(
@Value("${spring.kafka.bootstrap-servers}") String bootstrapServers,
@@ -39,7 +41,9 @@ public class KafkaConfig {
}
// 第二个Kafka配置
/**
* 分支中心kafka配置湖仓Kafka配置
*/
@Bean(name = "analyticCenterKafkaTemplate")
public KafkaTemplate<String, String> analyticCenterKafkaTemplate(
@Value("${analyticCenterKafka.bootstrap-servers}") String bootstrapServers,
@@ -54,6 +58,10 @@ public class KafkaConfig {
return new KafkaTemplate<>(factory);
}
/**
* 分支中心kafka消费者配置
*/
@Bean(name = "analyticCenterConsumerFactory")
public ConcurrentKafkaListenerContainerFactory<String, String> analyticCenterConsumerFactory(
@Value("${analyticCenterKafka.bootstrap-servers}") String bootstrapServers,
@@ -75,7 +83,6 @@ public class KafkaConfig {
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(props));
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
// factory.setBatchListener(true);
return factory;
}

View File

@@ -1,6 +1,5 @@
package com.volvo.ai.analytic.center.controller;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.ai.analytic.center.service.TmNameplateCorpusService;
import com.volvo.common.core.util.ResultMsg;
import io.swagger.annotations.Api;
@@ -25,10 +24,8 @@ public class NameplateCorpusController {
@Autowired
private TmNameplateCorpusService tmNameplateCorpusService;
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@PostMapping("/update")
@ApiOperation(value = "更新dify结果")
@ApiOperation(value = "dify发起http 东西网关调用,异步更新dify结果")
public ResultMsg<Object> updateNameplate(@RequestBody String message) {
log.info("updateNameplate message: {}", message);
return tmNameplateCorpusService.updateNameplate(message);

View File

@@ -3,9 +3,6 @@ package com.volvo.ai.analytic.center.controller;
import com.alibaba.fastjson.JSONObject;
import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO;
import com.volvo.ai.analytic.center.entity.TmTelephoneCorpus;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
import com.volvo.common.core.util.ResultMsg;
import io.swagger.annotations.Api;
import io.swagger.annotations.ApiOperation;
@@ -16,9 +13,10 @@ 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.*;
import javax.sql.DataSource;
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;
@RestController
@@ -31,9 +29,6 @@ public class TestController {
@Autowired
private RocketMQTemplate rocketMQTemplate;
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@Autowired
@Qualifier("analyticCenterKafkaTemplate")
private KafkaTemplate<String, String> kafkaProducer;
@@ -42,9 +37,6 @@ 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) {
@@ -70,26 +62,15 @@ public class TestController {
return ResultMsg.ok("ok");
}
@PostMapping("/mockDccKafka")
@ApiOperation(value = "补偿处理消息")
@ApiOperation(value = "手动发送kafka消息用于开发测试")
public ResultMsg<Object> mockDccKafka(@RequestBody AicorpusTelephoneDTO message) {
dccKafkaProducer.send("topic_voc_covert_text_log",JSONObject.toJSONString(message));
for(int i=0;i<20;i++){
message.setSourceId("000001"+i);
dccKafkaProducer.send("topic_voc_covert_text_log",JSONObject.toJSONString(message));
}
return ResultMsg.ok("ok");
}
@GetMapping("/pool")
public String checkPool() {
return dataSource.getClass().getName();
}
@Autowired
public TmTelephoneCorpusService tmTelephoneCorpusService;
@PostMapping("/save")
public void save(@RequestBody TmTelephoneCorpus tmTelephoneCorpus) {
tmTelephoneCorpusService.saveTelephoneCorpus(tmTelephoneCorpus);
}
}

View File

@@ -10,13 +10,13 @@ import com.volvo.ai.analytic.center.dto.corpus.CorpusReportDTO;
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
import com.volvo.ai.analytic.center.entity.AiAnalysisRequestLogs;
import com.volvo.ai.analytic.center.entity.TmNameplateCorpus;
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
import com.volvo.ai.analytic.center.enums.CategoryEnum;
import com.volvo.ai.analytic.center.feign.DiFyFeign;
import com.volvo.ai.analytic.center.mapper.AiAnalysisRequestLogsMapper;
import com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper;
import com.volvo.ai.analytic.center.service.*;
import com.volvo.ai.analytic.center.utils.FlowResultSplitUtil;
import com.xxl.job.core.handler.annotation.XxlJob;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
@@ -78,6 +78,7 @@ public class CorpusFailJob {
private TmTelephoneCorpusMapper tmTelephoneCorpusMapper;
@Autowired
private TmNameplateCorpusService tmNameplateCorpusService;
private DiFyFeign diFyFeign;
@@ -89,14 +90,14 @@ public class CorpusFailJob {
public void corpusFailTask() {
log.info(" corpusFailTask解析失败重试处理");
Integer total = aiAnalysisErrorsService.queryCountAnalysisErrorList(Arrays.asList(BusinessTypeEnum.SMART_ASSISTANT.getCode(),BusinessTypeEnum.SMART_ASSISTANT_QIWEI.getCode()));
Integer total = aiAnalysisErrorsService.queryCountAnalysisErrorList(Arrays.asList(BusinessTypeEnum.SMART_ASSISTANT.getCode(),BusinessTypeEnum.SMART_ASSISTANT_QIWEI.getCode(),BusinessTypeEnum.SMART_ASSISTANT_NAMEPLATE.getCode()));
log.info("corpusFailTask语料解析失败重试处理数据量{}", total);
int totalPages = PageDto.getTotalPages(total, pageSize);
log.info("corpusFailTask 语料解析失败重试处理数据量:{},总页数:{}", total, totalPages);
for (int i = 1; i <= totalPages; i++) {
int offset = (i - 1) * pageSize;
List<AiAnalysisErrors> aiAnalysisErrorsListlist = aiAnalysisErrorsService.queryAnalysisErrorList(Arrays.asList(BusinessTypeEnum.SMART_ASSISTANT.getCode(),BusinessTypeEnum.SMART_ASSISTANT_QIWEI.getCode()),offset, pageSize);
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 -> {
@@ -107,38 +108,23 @@ public class CorpusFailJob {
queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisErrors.getAiAnalysisRequestId());
AiAnalysisRequestLogs oldAiAnalysisRequestLogs = aiAnalysisRequestLogsMapper.selectOne(queryWrapper);
if (null != oldAiAnalysisRequestLogs) {
if (null != oldAiAnalysisRequestLogs && StringUtils.isBlank(oldAiAnalysisRequestLogs.getBusinessResponse())) {
DiFyReq diFyReq = JSONObject.parseObject(oldAiAnalysisRequestLogs.getDifyRequest(), DiFyReq.class);
CorpusReportDTO corpusReportDTO = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class);
Map<String, String> ltoMap = new HashMap<>();
ltoMap.put("analysisRecordId", oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
ltoMap.put("analysisScene", corpusReportDTO.getAnalysisScene() + "");
String tag = "";
if (null != corpusReportDTO && corpusReportDTO.getAnalysisScene() == 1) { //企微
ltoMap.put("unionId", corpusReportDTO.getUnionId());
ltoMap.put("consultantId", corpusReportDTO.getUserId());
tag = CategoryEnum.ENTERPRISE_WECHAT.getCode();
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyReq);
log.info(" corpusFailTask runDify execDifyFlow {}", execDifyFlow);
if (null != execDifyFlow && execDifyFlow.get("status").equals("succeeded")) {
String text = execDifyFlow.getJSONObject("outputs").getString("text");
String resultStrOne = FlowResultSplitUtil.flowOutputTextSplit(text, "任务1", "任务2");
String resultStrTwo = FlowResultSplitUtil.flowOutputTextSplit(text, "任务2", null);
if (StringUtils.isBlank(resultStrOne) || StringUtils.isBlank(resultStrTwo)) {
log.info(" corpusFailTask 企微语料解析为空text:{}", text);
return;
}
ltoMap.put("communicateDate", corpusReportDTO.getCorpusTime());
ltoMap.put("analysisResult", resultStrOne.replace("#", ""));
ltoMap.put("analysisDetail", resultStrTwo);
JSONObject text = execDifyFlow.getJSONObject("outputs");
// 发送MQ
log.info("corpusFailTask send mq {}", ltoMap);
tmTelephoneCorpusService.sendMq(tag, JSONObject.toJSONString(ltoMap));
log.info("send mq {}", text);
tmTelephoneCorpusService.sendMq(tag, text.toJSONString());
// 保存报告
oldAiAnalysisRequestLogs.setBusinessResponse(JSONObject.toJSONString(ltoMap));
oldAiAnalysisRequestLogs.setBusinessResponse( text.toJSONString());
oldAiAnalysisRequestLogs.setDifyResponse(JSON.toJSONString(execDifyFlow));
updateDiFyRequest(oldAiAnalysisRequestLogs, oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
@@ -147,7 +133,15 @@ public class CorpusFailJob {
updateAiAnalysisErrors(aiAnalysisErrors, oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
}
} else {
} else if(null != corpusReportDTO && corpusReportDTO.getAnalysisScene() == 3){
log.info(" corpusFailTask 铭牌数据重试:{},aiId:{}", corpusReportDTO.getCustomerFlowId(), oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyReq, oldAiAnalysisRequestLogs.getAiAnalysisRequestType(), JSONObject.toJSONString(corpusReportDTO),oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
List<TmNameplateCorpus> nameplateCorpusList = tmNameplateCorpusService.queryTelephoneCorpusByCustomerFlowId( Arrays.asList(corpusReportDTO.getCustomerFlowId()));
if(CollectionUtils.isNotEmpty(nameplateCorpusList)) {
TmNameplateCorpus nameplateCorpus = nameplateCorpusList.get(0);
tmNameplateCorpusService.sendNameplateLto(execDifyFlow,oldAiAnalysisRequestLogs.getAiAnalysisRequestId(), nameplateCorpus);
}
}else {
List<AicorpusTelephoneDTO> dccDtoList = tmTelephoneCorpusMapper.queryTelephoneCorpusBySourceIds( Arrays.asList(corpusReportDTO.getRecordId()));
if(CollectionUtils.isNotEmpty(dccDtoList)){
AicorpusTelephoneDTO dccDto = dccDtoList.get(0);

View File

@@ -1,15 +1,12 @@
package com.volvo.ai.analytic.center.job;
import com.volvo.ai.analytic.center.mapper.TmTelephoneCorpusMapper;
import com.volvo.ai.analytic.center.service.TmNameplateCorpusService;
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
import com.volvo.common.core.util.ResultMsg;
import com.xxl.job.core.context.XxlJobHelper;
import com.xxl.job.core.handler.annotation.XxlJob;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.PostMapping;
@@ -29,16 +26,10 @@ public class NameplateCorpusJob {
@Autowired
private TmNameplateCorpusService tmNameplateCorpusService;
@Autowired
private TmTelephoneCorpusService tmTelephoneCorpusService;
@Autowired
private TmTelephoneCorpusMapper tmTelephoneCorpusMapper;
@Value("${dify.corpus.nameplate.isUse}")
private boolean isUse;
/**
* 铭牌语料处理失败重试
* @param paramJson
* 手动铭牌语料处理失败重试
* @param paramJson statTimeendTimecustomerFlowIds
* @return
*/
@XxlJob("nameplateCorpusTaskFailRetry")

View File

@@ -24,7 +24,7 @@ import org.springframework.web.client.RestTemplate;
/**
* @ClassName AnalysisDifyMqConsumer
* @Description AI解析MQ-Callback处理
* @Description AI解析MQ-Callback处理 (暂时未用)
* @Author renzhen
* @Date 2025-03-04 10:18
* @Version 1.0
@@ -44,17 +44,9 @@ public class AnalysisDifyCallbackMqConsumer implements RocketMQListener<MessageE
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@Value("${dify.corpus.checkDccRepeat}")
private String checkDccRepeat;
@Autowired
private DiFyService diFyService;
@Autowired
private RestTemplate restTemplate;
private final ObjectMapper objectMapper = new ObjectMapper();
@Override
public void onMessage(MessageExt messageExt) {
long startTime = System.currentTimeMillis();

View File

@@ -5,10 +5,8 @@ import com.alibaba.fastjson.JSONObject;
import com.volvo.ai.analytic.center.dto.req.DiFyReq;
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
import com.volvo.ai.analytic.center.service.AiAnalysisErrorsService;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.ai.analytic.center.service.DiFyService;
import com.volvo.ai.analytic.center.utils.ConstantStr;
import com.volvo.ai.analytic.center.utils.RedisCounterRateLimiter;
import com.volvo.ai.analytic.center.utils.RedisLockService;
import com.volvo.ai.analytic.center.utils.RedisZSetUtil;
import lombok.extern.slf4j.Slf4j;
@@ -18,7 +16,6 @@ import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.data.redis.core.RedisTemplate;
import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.RestController;
@@ -27,7 +24,7 @@ import java.util.concurrent.CompletableFuture;
/**
* @ClassName AnalysisDifyMqConsumer
* @Description AI解析 MQ处理
* @Description AI解析 MQ处理 (暂时未用)
* @Author renzhen
* @Date 2025-03-04 10:18
* @Version 1.0
@@ -44,12 +41,6 @@ import java.util.concurrent.CompletableFuture;
public class AnalysisDifyMqConsumer implements RocketMQListener<MessageExt> {
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
@Value("${dify.corpus.checkDccRepeat}")
private String checkDccRepeat;
@Value("${rocketmq.consumer.analysisDify.difyLimit}")
private int difyLimit;
@@ -58,14 +49,9 @@ public class AnalysisDifyMqConsumer implements RocketMQListener<MessageExt> {
@Autowired
private DiFyService diFyService;
@Autowired
private RedisCounterRateLimiter redisCounterRateLimiter;
@Autowired
private AiAnalysisErrorsService aiAnalysisErrorsService;
@Autowired
private RedisTemplate<String, String> redisTemplate;
@Autowired
private RedisZSetUtil redisZSetUtil;

View File

@@ -4,7 +4,6 @@ import com.alibaba.fastjson.JSONObject;
import com.volvo.ai.analytic.center.entity.TmAnalysisResult;
import com.volvo.ai.analytic.center.entity.TtAnalysisResultInfo;
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
import com.volvo.ai.analytic.center.mapper.AiAnalysisRequestLogsMapper;
import com.volvo.ai.analytic.center.mapper.TtAnalysisResultInfoMapper;
import com.volvo.ai.analytic.center.service.TmAnalysisResultService;
import lombok.extern.slf4j.Slf4j;
@@ -19,6 +18,9 @@ import java.util.Arrays;
import java.util.Date;
import java.util.List;
/**
* LTO跟进记录AI总结的点赞点踩
*/
@Slf4j
@Component
@RocketMQMessageListener(consumerGroup = "${rocketmq.consumer.corpus.isLikeTpoicGroup}",
@@ -29,9 +31,6 @@ public class CorpushIsLikeConsumer implements RocketMQListener<MessageExt>{
@Autowired
private TmAnalysisResultService tmAnalysisResultService;
@Autowired
private AiAnalysisRequestLogsMapper aiAnalysisRequestLogsMapper;
@Autowired
private TtAnalysisResultInfoMapper ttAnalysisResultInfoMapper;
@Override

View File

@@ -4,7 +4,6 @@ package com.volvo.ai.analytic.center.mq;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.volvo.ai.analytic.center.dto.req.FourInOneRequestDTO;
import com.volvo.ai.analytic.center.service.AiAnalysisRequestLogsService;
import com.volvo.ai.analytic.center.service.IntelligentCustomerService;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.common.message.MessageExt;
@@ -17,7 +16,7 @@ import org.springframework.web.bind.annotation.RestController;
/**
* @ClassName CorpusProcessKafkaConsumer
* @Description 智能客服 语料解析 四合一
* @Description 消费4In1智能工单-在线客服语料消息
* @Author renzhen
* @Date 2025-03-04 10:18
* @Version 1.0
@@ -37,8 +36,6 @@ public class IntelligentCustomerMqConsumer implements RocketMQListener<MessageEx
@Autowired
private IntelligentCustomerService intelligentCustomerService;
@Autowired
private AiAnalysisRequestLogsService aiAnalysisRequestLogsService;
private final ObjectMapper objectMapper = new ObjectMapper();
@Override

View File

@@ -3,28 +3,22 @@ package com.volvo.ai.analytic.center.mq;
import com.alibaba.fastjson.JSON;
import com.volvo.ai.analytic.center.dto.corpus.NameplateTableKafkaDTO;
import com.volvo.ai.analytic.center.service.CorpusPortraitService;
import com.volvo.ai.analytic.center.service.TmNameplateCorpusService;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
import org.apache.commons.lang3.StringUtils;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource;
import java.util.concurrent.CompletableFuture;
/**
* @ClassName NameplateKafkaConsumer
* @Description 铭牌数据处理
* @Description 消费tm_nameplate_corpus表binlog的Kafka消息
* @Author renzhen
* @Date 2025-03-04 10:18
* @Version 1.0
@@ -40,24 +34,13 @@ public class NameplateKafkaConsumer {
private TmNameplateCorpusService tmNameplateCorpusService;
@Resource
private RocketMQTemplate rocketMqTemplate;
@Autowired
@Resource(name = "threadPoolTaskExecutor")
private ThreadPoolTaskExecutor executor;
@Autowired
private CorpusPortraitService corpusPortraitService;
@KafkaListener(topics = "${analyticCenterKafka.consumer.topic}",
groupId = "${analyticCenterKafka.consumer.group}" ,
containerFactory = "analyticCenterConsumerFactory",
concurrency = "3")
public void listen(String recordMessages, Acknowledgment ack,
@Header(KafkaHeaders.RECEIVED_PARTITION_ID) Integer partitionId,
@Header(KafkaHeaders.OFFSET) Long offset) throws InterruptedException {
@Header(KafkaHeaders.OFFSET) Long offset) {
long startTime = System.currentTimeMillis();
log.info("nameplateKafkaConsumer 当前线程: {}, 线程ID: {},计数:{}", Thread.currentThread().getName(), Thread.currentThread().getId());
log.info("nameplateKafkaConsumerMessage: {}", recordMessages);
@@ -73,19 +56,7 @@ public class NameplateKafkaConsumer {
return;
}
tmNameplateCorpus.getData().forEach(nameplate -> {
tmNameplateCorpusService.processItem(nameplate);
CompletableFuture.runAsync(() -> {
try {
corpusPortraitService.portraitNameplate(nameplate);
} catch (Exception e) {
log.error("corpusPortrait画像铭牌异步任务执行失败", e);
}
}, executor);
}
);
tmNameplateCorpus.getData().forEach(nameplate -> tmNameplateCorpusService.processItem(nameplate));
log.info("nameplateKafkaConsumeracknowledge{}",partitionId, offset);
// 手动提交 offset
} catch (Exception e) {

View File

@@ -10,10 +10,8 @@ import com.volvo.ai.analytic.center.dto.req.RunMaskingRuleInput;
import com.volvo.ai.analytic.center.entity.*;
import com.volvo.ai.analytic.center.enums.BusinessTypeEnum;
import com.volvo.ai.analytic.center.enums.CategoryEnum;
import com.volvo.ai.analytic.center.feign.RemoteCarModelClient;
import com.volvo.ai.analytic.center.mapper.AiAnalysisErrorsMapper;
import com.volvo.ai.analytic.center.mapper.TmNameplateCorpusMapper;
import com.volvo.ai.analytic.center.mapper.TmOdsVdqwMessagearchivingMapper;
import com.volvo.ai.analytic.center.mapper.TtNameplateRecordMapper;
import com.volvo.ai.analytic.center.service.*;
import com.volvo.ai.analytic.center.utils.ConstantStr;
@@ -21,17 +19,16 @@ import com.volvo.common.core.util.ResultMsg;
import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils;
import org.apache.commons.lang3.StringUtils;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.time.LocalDate;
import java.util.*;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
/**
@@ -44,8 +41,6 @@ import java.util.concurrent.CompletableFuture;
@Service
public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusMapper, TmNameplateCorpus> implements TmNameplateCorpusService {
@Autowired
private TmOdsVdqwMessagearchivingMapper tmOdsVdqwMessagearchivingMapper;
@Autowired
private TmNameplateCorpusMapper tmNameplateCorpusMapper;
@@ -62,27 +57,14 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
private AiAnalysisErrorsMapper aiAnalysisErrorsMapper;
@Autowired
private DiFyService diFyService;
@Resource
private RocketMQTemplate rocketMqTemplate;
@Value("${dify.corpus.nameplate.appkey}")
private String nameplateAppKey;
@Value("${batch.size}")
public int pageSize = 100;
@Autowired
private RemoteCarModelClient remoteCarModelClient;
@Autowired
private DataMaskingRuleService dataMaskingRuleService;
@Autowired
@Resource(name = "threadPoolTaskExecutor")
private ThreadPoolTaskExecutor executor;
@Autowired
private CorpusPortraitService corpusPortraitService;
@Override
public void runNameplateCorpusDifyRetry(String paramJson) {
long startTime = System.currentTimeMillis();
@@ -110,22 +92,33 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
Integer total = tmNameplateCorpusMapper.countQueryTmNameplateCorpusRetry(statTime, endTime, customerFlowIds ,retry);
int totalPages = PageDto.getTotalPages(total, pageSize);
// 获取消息列表
int optimalThreadPoolSize = Runtime.getRuntime().availableProcessors() + 1;
log.info("获取的线程数:{}",optimalThreadPoolSize);
// 创建线程池
ExecutorService executor = Executors.newFixedThreadPool(optimalThreadPoolSize); // 根据需求调整线程池大小
for (int i = 1; i <= totalPages; i++) {
int offset = (i - 1) * pageSize;
List<TmNameplateCorpus> messageList = tmNameplateCorpusMapper.queryTmNameplateCorpusRetry(statTime, endTime, offset, pageSize, customerFlowIds, retry);
for(TmNameplateCorpus tmNameplateCorpus :messageList){
CompletableFuture.runAsync(() -> {
try {
processItem(tmNameplateCorpus);
// 增加 画像手动补偿
corpusPortraitService.portraitNameplate(tmNameplateCorpus);
} catch (Exception e) {
log.error("重跑铭牌语料失败: customerFlowId={}, AcceptUserId={}, 异常: {}",
tmNameplateCorpus.getCustomerFlowId(), e.getMessage(), e);
}
}, executor);
}
// 处理查询到的数据
// 使用 CompletableFuture 并行处理
CompletableFuture<?>[] futures = messageList.stream()
.map(item -> CompletableFuture.runAsync(() -> {
try {
processItem(item);
} catch (Exception e) {
log.error("重跑铭牌语料失败: customerFlowId={}, AcceptUserId={}, 异常: {}",
item.getCustomerFlowId(), e.getMessage(), e);
}
}, executor))
.toArray(CompletableFuture[]::new);
// 等待所有任务完成
CompletableFuture.allOf(futures).join();
}
// 关闭线程池
executor.shutdown();
log.info("重跑铭牌语料铭牌数据跑批结束 耗时:{}",System.currentTimeMillis()-startTime);
}
@@ -230,9 +223,11 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
if(StringUtils.isNotEmpty(data)){
List<TmNameplateCorpus> tmNameplateCorpusList = new ArrayList<>();
TmNameplateCorpus analysisResp = JSONObject.parseObject(data, TmNameplateCorpus.class);
analysisResp.setCustomerFlowId(analysisResp.getCustomerFlowId());
tmNameplateCorpusList.add(analysisResp);
for(int i=0;i<100;i++){
TmNameplateCorpus analysisResp = JSONObject.parseObject(data, TmNameplateCorpus.class);
analysisResp.setCustomerFlowId(analysisResp.getCustomerFlowId()+i);
tmNameplateCorpusList.add(analysisResp);
}
this.saveBatch(tmNameplateCorpusList);
return ResultMsg.ok();
}

View File

@@ -334,10 +334,4 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
}
}, 10000);
}
}