修改kafka偏移量

This commit is contained in:
zren25
2025-03-31 19:21:45 +08:00
parent ed185efc67
commit a69a3757eb
6 changed files with 119 additions and 67 deletions

View File

@@ -3,6 +3,7 @@ package com.volvo.ai.analytic.center.job;
import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject; import com.alibaba.fastjson.JSONObject;
import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper;
import com.volvo.ai.analytic.center.dto.PageDto;
import com.volvo.ai.analytic.center.dto.corpus.CorpusReportDTO; import com.volvo.ai.analytic.center.dto.corpus.CorpusReportDTO;
import com.volvo.ai.analytic.center.dto.req.DiFyReq; import com.volvo.ai.analytic.center.dto.req.DiFyReq;
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors; import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
@@ -53,7 +54,8 @@ public class CorpusFailJob {
@Resource @Resource
private RocketMQTemplate rocketMqTemplate; private RocketMQTemplate rocketMqTemplate;
@Value("${batch.size}")
public int pageSize = 100;
@Autowired @Autowired
private TmTelephoneCorpusService tmTelephoneCorpusService; private TmTelephoneCorpusService tmTelephoneCorpusService;
@@ -65,69 +67,78 @@ public class CorpusFailJob {
public void corpusFailTask() { public void corpusFailTask() {
log.info("语料解析失败重试处理"); log.info("语料解析失败重试处理");
List<AiAnalysisErrors> aiAnalysisErrorsListlist = aiAnalysisErrorsService.queryAnalysisErrorList(BusinessTypeEnum.SMART_ASSISTANT.getCode()); Integer total = aiAnalysisErrorsService.queryCountAnalysisErrorList(BusinessTypeEnum.SMART_ASSISTANT.getCode());
if(CollectionUtils.isNotEmpty(aiAnalysisErrorsListlist)) { log.info("语料解析失败重试处理数据量:{}", total);
log.info("语料解析失败重试处理 size:{}", aiAnalysisErrorsListlist.size()); int totalPages = PageDto.getTotalPages(total, pageSize);
aiAnalysisErrorsListlist.stream().forEach(aiAnalysisErrors -> { log.info("语料解析失败重试处理数据量:{},总页数:{}", total, totalPages);
for (int i = 1; i <= totalPages; i++) {
int offset = (i - 1) * pageSize;
List<AiAnalysisErrors> aiAnalysisErrorsListlist = aiAnalysisErrorsService.queryAnalysisErrorList(BusinessTypeEnum.SMART_ASSISTANT.getCode(),offset, pageSize);
if(CollectionUtils.isNotEmpty(aiAnalysisErrorsListlist)) {
log.info("语料解析失败重试处理 size:{}", aiAnalysisErrorsListlist.size());
aiAnalysisErrorsListlist.stream().forEach(aiAnalysisErrors -> {
try { try {
LambdaQueryWrapper<AiAnalysisRequestLogs> queryWrapper = new LambdaQueryWrapper<>(); LambdaQueryWrapper<AiAnalysisRequestLogs> queryWrapper = new LambdaQueryWrapper<>();
queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisErrors.getAiAnalysisRequestId()); queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisErrors.getAiAnalysisRequestId());
AiAnalysisRequestLogs oldAiAnalysisRequestLogs = aiAnalysisRequestLogsMapper.selectOne(queryWrapper); AiAnalysisRequestLogs oldAiAnalysisRequestLogs = aiAnalysisRequestLogsMapper.selectOne(queryWrapper);
if (null != oldAiAnalysisRequestLogs) { if (null != oldAiAnalysisRequestLogs) {
DiFyReq diFyReq = JSONObject.parseObject(oldAiAnalysisRequestLogs.getDifyRequest(), DiFyReq.class); DiFyReq diFyReq = JSONObject.parseObject(oldAiAnalysisRequestLogs.getDifyRequest(), DiFyReq.class);
CorpusReportDTO corpusReportDTO = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class); CorpusReportDTO corpusReportDTO = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class);
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;
}
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();
} else {
ltoMap.put("recordId", corpusReportDTO.getUserId());
tag = CategoryEnum.PHONE_VOICE.getCode();
}
ltoMap.put("communicateDate", corpusReportDTO.getCorpusTime());
ltoMap.put("analysisResult", resultStrOne.replace("#", ""));
ltoMap.put("analysisDetail", resultStrTwo);
// 发送MQ
log.info("corpusFailTask send mq {}", ltoMap);
tmTelephoneCorpusService.sendMq(tag, JSONObject.toJSONString(ltoMap));
// 保存报告
oldAiAnalysisRequestLogs.setBusinessResponse(JSONObject.toJSONString(ltoMap));
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;
} }
Map<String, String> ltoMap = new HashMap<>(); oldAiAnalysisRequestLogs.setDifyResponse(JSON.toJSONString(execDifyFlow));
ltoMap.put("analysisRecordId", oldAiAnalysisRequestLogs.getAiAnalysisRequestId()); updateDiFyRequest(oldAiAnalysisRequestLogs, oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
ltoMap.put("analysisScene", corpusReportDTO.getAnalysisScene() + ""); aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount() + 1);
String tag = ""; aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("1");
if (null != corpusReportDTO && corpusReportDTO.getAnalysisScene() == 1) { //企微 updateAiAnalysisErrors(aiAnalysisErrors, oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
ltoMap.put("unionId", corpusReportDTO.getUnionId());
ltoMap.put("consultantId", corpusReportDTO.getUserId());
tag = CategoryEnum.ENTERPRISE_WECHAT.getCode();
} else {
ltoMap.put("recordId", corpusReportDTO.getUserId());
tag = CategoryEnum.PHONE_VOICE.getCode();
}
ltoMap.put("communicateDate", corpusReportDTO.getCorpusTime());
ltoMap.put("analysisResult", resultStrOne.replace("#", ""));
ltoMap.put("analysisDetail", resultStrTwo);
// 发送MQ
log.info("corpusFailTask send mq {}", ltoMap);
tmTelephoneCorpusService.sendMq(tag, JSONObject.toJSONString(ltoMap));
// 保存报告
oldAiAnalysisRequestLogs.setBusinessResponse(JSONObject.toJSONString(ltoMap));
} }
oldAiAnalysisRequestLogs.setDifyResponse(JSON.toJSONString(execDifyFlow)); } catch (Exception e) {
updateDiFyRequest(oldAiAnalysisRequestLogs,oldAiAnalysisRequestLogs.getAiAnalysisRequestId()); log.info("语料解析失败补偿异常:{}", e);
aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("0");
aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount() + 1); aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount() + 1);
aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("1"); updateAiAnalysisErrors(aiAnalysisErrors, aiAnalysisErrors.getAiAnalysisRequestId());
updateAiAnalysisErrors(aiAnalysisErrors,oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
} }
} catch (Exception e) { });
log.info("语料解析失败补偿异常:{}",e); }
aiAnalysisErrors.setAiAnalysisErrorHandlingStatus("0"); }
aiAnalysisErrors.setRetryCount(aiAnalysisErrors.getRetryCount()+ 1);
updateAiAnalysisErrors(aiAnalysisErrors,aiAnalysisErrors.getAiAnalysisRequestId());
}
});
}
} }

View File

@@ -10,7 +10,7 @@ import java.util.List;
@Mapper @Mapper
public interface AiAnalysisErrorsMapper extends BaseMapper<AiAnalysisErrors> { public interface AiAnalysisErrorsMapper extends BaseMapper<AiAnalysisErrors> {
List<AiAnalysisErrors> queryAnalysisErrorList(@Param("businessType") String businessType);
int queryCountAnalysisErrorList(@Param("businessType") String businessType);
List<AiAnalysisErrors> queryAnalysisErrorList( @Param("businessType") String businessType, @Param("offset") int offset, @Param("pageSize") int pageSize);
} }

View File

@@ -17,6 +17,8 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.cloud.context.config.annotation.RefreshScope; import org.springframework.cloud.context.config.annotation.RefreshScope;
import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.annotation.PartitionOffset;
import org.springframework.kafka.annotation.TopicPartition;
import org.springframework.messaging.support.MessageBuilder; import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.PostMapping;
@@ -54,7 +56,27 @@ public class CorpusProcessKafkaProducer {
private RocketMQTemplate rocketMqTemplate; private RocketMQTemplate rocketMqTemplate;
@PostMapping("corpusProcessKafkaConsumer") @PostMapping("corpusProcessKafkaConsumer")
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}") @KafkaListener(
topicPartitions = @TopicPartition(
topic = "${spring.kafka.topic}",
partitions = {"0", "1","2", "3","4", "5","6", "7","8", "9", "10", "11"},
partitionOffsets = {
@PartitionOffset(partition = "0", initialOffset = "2792520"),
@PartitionOffset(partition = "1", initialOffset = "2596153"),
@PartitionOffset(partition = "2", initialOffset = "2536889"),
@PartitionOffset(partition = "3", initialOffset = "2782173"),
@PartitionOffset(partition = "4", initialOffset = "2616677"),
@PartitionOffset(partition = "5", initialOffset = "2529585"),
@PartitionOffset(partition = "6", initialOffset = "2761950"),
@PartitionOffset(partition = "7", initialOffset = "2581957"),
@PartitionOffset(partition = "8", initialOffset = "2530225"),
@PartitionOffset(partition = "9", initialOffset = "2803179"),
@PartitionOffset(partition = "10", initialOffset = "2599129"),
@PartitionOffset(partition = "11", initialOffset = "2546277")
}
),
groupId = "${spring.kafka.group}"
)
public void listen(List<String> recordMessages) { public void listen(List<String> recordMessages) {
long startTime = System.currentTimeMillis(); long startTime = System.currentTimeMillis();
try { try {
@@ -97,11 +119,6 @@ public class CorpusProcessKafkaProducer {
} }
}, 10000); }, 10000);
// CompletableFuture.runAsync(() -> {
// tmTelephoneCorpusService.runTelephoneCorpusDify(aicorpusTelephone);
// }, runTelephoneExecutor);
// // 关闭线程池
// runTelephoneExecutor.shutdown();
log.info(" dify处理耗时{}", System.currentTimeMillis() - startTimeDify); log.info(" dify处理耗时{}", System.currentTimeMillis() - startTimeDify);
} }
} catch (Exception e) { } catch (Exception e) {

View File

@@ -2,12 +2,13 @@ package com.volvo.ai.analytic.center.service;
import com.baomidou.mybatisplus.extension.service.IService; import com.baomidou.mybatisplus.extension.service.IService;
import com.volvo.ai.analytic.center.entity.AiAnalysisErrors; import com.volvo.ai.analytic.center.entity.AiAnalysisErrors;
import org.apache.ibatis.annotations.Param;
import java.util.List; import java.util.List;
public interface AiAnalysisErrorsService extends IService<AiAnalysisErrors> { public interface AiAnalysisErrorsService extends IService<AiAnalysisErrors> {
boolean saveAiAnalysisErrors(AiAnalysisErrors entity); boolean saveAiAnalysisErrors(AiAnalysisErrors entity);
int queryCountAnalysisErrorList(String businessType);
List<AiAnalysisErrors> queryAnalysisErrorList(String businessType); List<AiAnalysisErrors> queryAnalysisErrorList( String businessType, int offset,int pageSize);
} }

View File

@@ -30,9 +30,14 @@ public class AiAnalysisErrorsServiceImpl extends ServiceImpl<AiAnalysisErrorsMa
} }
} }
public List<AiAnalysisErrors> queryAnalysisErrorList(String businessType) { @Override
public int queryCountAnalysisErrorList(String businessType) {
return aiAnalysisErrorsMapper.queryCountAnalysisErrorList(businessType);
}
public List<AiAnalysisErrors> queryAnalysisErrorList( String businessType, int offset,int pageSize) {
//捞取异常表中属于社区的异常数据 //捞取异常表中属于社区的异常数据
return aiAnalysisErrorsMapper.queryAnalysisErrorList(businessType); return aiAnalysisErrorsMapper.queryAnalysisErrorList(businessType, offset, pageSize);
} }
} }

View File

@@ -1,6 +1,23 @@
<?xml version="1.0" encoding="UTF-8"?> <?xml version="1.0" encoding="UTF-8"?>
<!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd"> <!DOCTYPE mapper PUBLIC "-//mybatis.org//DTD Mapper 3.0//EN" "http://mybatis.org/dtd/mybatis-3-mapper.dtd">
<mapper namespace="com.volvo.ai.analytic.center.mapper.AiAnalysisErrorsMapper"> <mapper namespace="com.volvo.ai.analytic.center.mapper.AiAnalysisErrorsMapper">
<select id="queryCountAnalysisErrorList" resultType="java.lang.Integer">
select count(1) from (
SELECT
ai_analysis_request_id as aiAnalysisRequestId,
ai_analysis_request_type as aiAnalysisRequestType,
ai_analysis_error_handling_status as aiAnalysisErrorHandlingStatus,
retry_count as retryCount,
max_retry_count as maxRetryCount
FROM
tt_ai_analysis_errors
WHERE
is_deleted = 0
AND ai_analysis_request_type=#{businessType}
AND ai_analysis_error_handling_status = '0'
AND retry_count &lt;= max_retry_count
)tag
</select>
<select id="queryAnalysisErrorList" resultType="com.volvo.ai.analytic.center.entity.AiAnalysisErrors"> <select id="queryAnalysisErrorList" resultType="com.volvo.ai.analytic.center.entity.AiAnalysisErrors">
@@ -17,6 +34,7 @@
AND ai_analysis_request_type=#{businessType} AND ai_analysis_request_type=#{businessType}
AND ai_analysis_error_handling_status = '0' AND ai_analysis_error_handling_status = '0'
AND retry_count &lt;= max_retry_count AND retry_count &lt;= max_retry_count
LIMIT #{offset}, #{pageSize}
</select> </select>
</mapper> </mapper>