10 Commits

Author SHA1 Message Date
ZLI263
9f1397403f 添加日志,定位是否传 soureId 2025-09-24 17:20:09 +08:00
ZLI263
ae5035e5d0 去掉定位日志 2025-09-24 13:16:33 +08:00
ZLI263
2b04df3b7d 解析格式化时间 2025-09-24 12:41:32 +08:00
ZLI263
f727175bca 铭牌如果时间为空,处理为null 2025-09-24 12:24:57 +08:00
ZLI263
c9399d7c41 铭牌如果时间为空,处理为null 2025-09-24 12:07:33 +08:00
ZLI263
1cede9f729 铭牌加定位日志 2025-09-24 12:05:10 +08:00
ZLI263
f44fcd2bfc 铭牌加定位日志 2025-09-24 12:04:47 +08:00
ZLI263
9c5df51d39 铭牌加定位日志 2025-09-24 11:59:55 +08:00
ZLI263
fa91dd65d4 名片加定位日志 2025-09-24 11:37:16 +08:00
ZLI263
5ffd0e5f8f 解决线程池过度关闭问题 2025-09-24 10:33:24 +08:00
3 changed files with 55 additions and 35 deletions

View File

@@ -115,21 +115,17 @@ public class CorpusProcessKafkaProducer {
}
} catch (Exception e) {
log.error("CorpusProcessKafkaProducer 电话语料 解析JSON出错: {}", e.getMessage());
}finally {
executor.shutdown();
}
});
}
executor.shutdown();
}
log.info("Kafka 消息处理完成,耗时:{}", System.currentTimeMillis() - startTime);
// 在这里可以添加对解析后的对象的进一步处理逻辑
} catch (Exception e) {
log.error("CorpusProcessKafkaProducer 电话语料 解析JSON出错: {}" , e.getMessage());
}finally {
executor.shutdown();
}
}

View File

@@ -26,7 +26,10 @@ import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
import org.springframework.stereotype.Service;
import javax.annotation.Resource;
import java.text.SimpleDateFormat;
import java.time.LocalDate;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.time.ZonedDateTime;
import java.time.format.DateTimeFormatter;
import java.util.*;
@@ -145,7 +148,7 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
try {
Optional.ofNullable(item).filter(tmNameplateCorpus -> tmNameplateCorpus.getCustomerFlowId() != null && tmNameplateCorpus.getNameplateContent() != null).orElseThrow(() -> new RuntimeException("铭牌语料为空"));
log.info("processItem开始处理铭牌语料:{}", item);
String customerFlowId = item.getCustomerFlowId();
String nameplateContent = item.getNameplateContent();
log.info("铭牌数据处理:customerFlowId:{}", customerFlowId);
@@ -203,7 +206,6 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
long endTime = System.currentTimeMillis();
log.info("第一个业务场景(总结和分类)执行时间: {} ms", (endTime - startTime));
//第一个业务场景, 结束
long startTime2 = System.currentTimeMillis();
try {
log.info("铭牌语料可用许可授权数,画像场景={}", semaphore.availablePermits());
@@ -213,17 +215,33 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
log.warn("获取许可超时,跳过执行任务");
return;
}
log.info("itemvalue 画像业务场景{}",item);
CompletableFuture.runAsync(() -> {
ZonedDateTime zonedDateTime = ZonedDateTime.parse(item.getNameplateEndTime().toString());
DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
String formattedDateStartTime = zonedDateTime.format(formatter);
//"analysisScene": "分析类型", // 1、企微会话 2、AI通话录音 3、AI铭牌(客流) 4、AI铭牌(试驾)
DiFyReq diFyImageReq2 = new DiFyReq();
diFyImageReq2.setUser(ConstantStr.corpus_user);
Map<String, Object> inputMap2 = new HashMap<>();
if(item.getNameplateEndTime() == null){
log.info("getNameplateEndTime is null");
inputMap2.put("communicateDate", "null");
}else {
try {
log.info("getNameplateEndTime is not null");
// 解析为 LocalDateTime(假设 CST 是 Asia/Shanghai)
LocalDateTime ldt = LocalDateTime.parse(item.getNameplateEndTime().toString().replace(" CST ", " "),
DateTimeFormatter.ofPattern("EEE MMM dd HH:mm:ss yyyy", Locale.ENGLISH));
// 转为带时区的时间(可选)
java.time.ZonedDateTime zonedDateTime = ldt.atZone(ZoneId.of("Asia/Shanghai"));
// 格式化输出
DateTimeFormatter outputFormatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
String formattedDate = zonedDateTime.format(outputFormatter);
inputMap2.put("communicateDate", formattedDate);
}catch (Exception e){
log.info("getNameplateEndTime {}", e.getMessage());
}
}
diFyImageReq2.setUser(ConstantStr.corpus_user);
inputMap2.put("businessId", item.getCustomerFlowId());
inputMap2.put("communicateDate", formattedDateStartTime);
inputMap2.put("analysisScene", "3");
inputMap2.put("version", 2);
inputMap2.put("chat", corpusChat);
@@ -239,7 +257,7 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
}, executor);
} catch (Exception e) {
log.error("执行Dify失败: customerFlowId={}, 错误信息: {}",
log.error("DifyFailure: customerFlowId={}, 错误信息: {}",
item.getCustomerFlowId(), e.getMessage(), e);
} finally {
// 释放许可
@@ -342,7 +360,7 @@ public class TmNameplateCorpusServiceImpl extends ServiceImpl<TmNameplateCorpusM
aiAnalysisErrorsMapper.update(AiAnalysisErrors.builder().aiAnalysisRequestId(oldAiAnalysisErrors.getAiAnalysisRequestId()).aiAnalysisErrorHandlingStatus("1").build(), errorQueryWrapper);
}
} else {
log.info("没找到语料记录");
log.info("没找到语料记录,{}",message);
}
return ResultMsg.ok();
}

View File

@@ -116,13 +116,10 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
RunMaskingRuleInput runMaskingRuleInput = new RunMaskingRuleInput();
runMaskingRuleInput.setDataMaskingRules(maskingRuleItems);
Map<String, Object> inputMap = new HashMap();
DiFyReq diFyImageReq = new DiFyReq();
diFyImageReq.setUser(ConstantStr.corpus_user);
diFyImageReq.setFlowId(telephoneToken);
JSONObject jsonObject = JSONObject.parseObject(aicorpusTelephone.getDisplay());
JSONArray segments = jsonObject.getJSONArray("segments");
Long audioDuration = jsonObject.getLong("audio_duration"); // 毫秒
if (audioDuration / 1000 <= 10) {
log.info("电话语料时长小于10秒,不进行dify处理");
return;
@@ -157,14 +154,20 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
ZonedDateTime zonedDateTime = ZonedDateTime.parse(jsonObject.getString("start_time"));
DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
String formattedDateStartTime = zonedDateTime.format(formatter);
String businessIdSourceId= aicorpusTelephone.getSourceId();
String carModel = getCarModelList();
inputMap.put("chat", chatList.toString());
inputMap.put("model", carModel);
inputMap.put("analysisScene", "2");
inputMap.put("recordId", aicorpusTelephone.getSourceId());
inputMap.put(communicateDateStr, formattedDateStartTime);
inputMap.put("version", 2);
diFyImageReq.setInputs(inputMap);
Map<String, Object> inputMapForSummary = new HashMap<>();
inputMapForSummary.put("chat", chatList.toString());
inputMapForSummary.put("model", carModel);
inputMapForSummary.put("analysisScene", "2");
inputMapForSummary.put("recordId", aicorpusTelephone.getSourceId());
inputMapForSummary.put(communicateDateStr, formattedDateStartTime);
inputMapForSummary.put("version", 2);
Map<String, Object> inputMapForPortrait = new HashMap<>(inputMapForSummary); // 复制一份
inputMapForPortrait.put("businessId", businessIdSourceId); // 仅画像任务需要
CorpusReportDTO corpusReportDTO = new CorpusReportDTO();
corpusReportDTO.setCorpusTime(formattedDateStartTime);
corpusReportDTO.setRecordId(aicorpusTelephone.getSourceId());
@@ -175,9 +178,13 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
long startTime = System.currentTimeMillis();
CompletableFuture.runAsync(() -> {
try {
diFyImageReq.setFlowId(telephoneToken);
log.info("dcc语料telephoneToken:{}", telephoneToken);
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode(),
DiFyReq req1 = new DiFyReq();
req1.setUser(ConstantStr.corpus_user);
req1.setFlowId(telephoneToken);
req1.setInputs(inputMapForSummary);
log.info("dcc语料req1:{}", req1);
JSONObject execDifyFlow = diFyService.executeDifyFlow(req1, BusinessTypeEnum.SMART_ASSISTANT.getCode(),
JSONObject.toJSONString(corpusReportDTO), aicorpusTelephone.getAiAnalysisRequestId());
log.info("dcc总结场景,总结和分类 runDify execDifyFlow : {}", execDifyFlow);
@@ -194,15 +201,14 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
long startTime2 = System.currentTimeMillis();
CompletableFuture.runAsync(() -> {
try {
// 创建新的DiFyReq对象以避免线程安全问题
inputMap.put("businessId", aicorpusTelephone.getSourceId());
DiFyReq req2 = new DiFyReq();
req2.setUser(ConstantStr.corpus_user);
req2.setFlowId(dccTokenPortrait);
req2.setInputs(inputMapForPortrait); // 使用带 businessId 的副本
diFyImageReq.setInputs(inputMap);
diFyImageReq.setFlowId(dccTokenPortrait);
log.info("dcc总结场景,客户画像runDify execDifyFlow : {}", dccTokenPortrait);
JSONObject execDifyFlowForPortrait = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.CORPUS_PORTRAIT_DCC.getCode(),
JSONObject execDifyFlowForPortrait = diFyService.executeDifyFlow(req2, BusinessTypeEnum.CORPUS_PORTRAIT_DCC.getCode(),
JSONObject.toJSONString(corpusReportDTO), aicorpusTelephone.getAiAnalysisRequestId());
log.info("dcc客户画像场景 runDify execDifyFlow 返回 , dcc: {}", execDifyFlowForPortrait);