@@ -23,7 +23,6 @@ import com.volvo.ai.analytic.center.mq.CorpusQuestionProducer;
import com.volvo.ai.analytic.center.service.* ;
import com.volvo.ai.analytic.center.utils.AiAnalysisUtils ;
import com.volvo.ai.analytic.center.utils.ConstantStr ;
import com.volvo.ai.analytic.center.utils.FlowResultSplitUtil ;
import lombok.extern.slf4j.Slf4j ;
import org.apache.commons.collections.CollectionUtils ;
import org.apache.commons.lang3.StringUtils ;
@@ -107,151 +106,146 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
public void saveTelephoneCorpus ( TmTelephoneCorpus tmTelephoneCorpus ) {
this . save ( tmTelephoneCorpus ) ;
}
private String communicateDateStr = " communicateDate " ;
private String outputsStr = " outputs " ;
private String aiAnalysisRequestIdStr = " aiAnalysisRequestId " ;
@Override
public void runTelephoneCorpusDify ( AicorpusTelephoneDTO aicorpusTelephone ) {
try {
if ( null ! = aicorpusTelephone ) {
if ( null ! = aicorpusTelephone ) {
List < DataMaskingRule > maskingRuleItems = dataMaskingRuleService . getDataMaskingRuleListByApplicationChannel ( BusinessTypeEnum . SMART_ASSISTANT . getCode ( ) ) ;
List < Data MaskingRule> m askingRuleItems = dataMaskingRuleService . getDataMaskingRuleListByApplicationChannel ( BusinessTypeEnum . SMART_ASSISTANT . getCode ( ) ) ;
Run MaskingRuleInput runM askingRuleInput = new RunMaskingRuleInput ( ) ;
runMaskingRuleInput . setDataMaskingRules ( maskingRuleItems ) ;
RunMaskingRuleInput runMaskingRuleInput = new RunMaskingRuleInput ( ) ;
runMaskingRuleInput . setDataMaskingRules ( maskingRuleItems ) ;
Map < String , Object > inputMap = new HashMap ( ) ;
DiFyReq diFyImageReq = new DiFyReq ( ) ;
diFyImageReq . setUser ( ConstantStr . corpus_user ) ;
Map < String , Object > inputMap = new HashMap ( ) ;
DiFyReq diFyImageReq = new DiFyReq ( ) ;
diFyImageReq . setUser ( ConstantStr . corpus_user ) ;
JSONObject jsonObject = JSONObject . parseObject ( aicorpusTelephone . getDisplay ( ) ) ;
log . info ( " 电话语料内容:{} " , jsonObject . toString ( ) ) ;
inputMap . put ( " chat " , aicorpusTelephone . getDisplay ( ) ) ;
JSONArray segments = jsonObject . getJSONArray ( " segments " ) ;
log . info ( " 电话语料内容segments: {} " , segments . toString ( ) ) ;
// 遍历 segments
StringBuffer chatList = new StringBuffer ( ) ;
segments . stream ( )
. map ( segment - > ( JSONObject ) segment )
. forEach ( segment - > {
JSONObject result = segment . getJSONObject ( " result " ) ;
String text = result . getString ( " text " ) ;
JSONObject analysisInfo = result . getJSONObject ( " analysis_info " ) ;
String role = analysisInfo . getString ( " role " ) ;
String title = " " ;
if ( role . equals ( " AGENT " ) ) {
title = " 顾问 " ;
} else {
title = " 客户 " ;
}
// 拼接 role 和 text
String chat = title + " : " + text ;
runMaskingRuleInput . setOldStr ( chat ) ;
String corpusChat = dataMaskingRuleService . runMaskingRule ( runMaskingRuleInput ) ;
chatList . append ( corpusChat ) . append ( " \ n " ) ;
JSONObject jsonObject = JSONObject . parseObject ( aicorpusTelephone . getDisplay ( ) ) ;
} ) ;
inputMap . put ( " chat " , aicorpusTelephone . getDisplay ( ) ) ;
JSONArray segments = jsonObject . getJSONArray ( " segment s" ) ;
ZonedDateTime zonedDateTime = ZonedDateTime . parse ( jsonObject . getString ( " start_time " ) ) ;
DateTimeFormatter formatter = DateTimeFormatter . ofPattern ( " yyyy-MM-dd HH:mm:s s" ) ;
String formattedDateStartTime = zonedDateTime . format ( formatter ) ;
// 遍历 segments
StringBuffer chatList = new StringBuffer ( ) ;
segments . stream ( )
. map ( segment - > ( JSONObject ) segment )
. forEach ( segment - > {
JSONObject result = segment . getJSONObject ( " result " ) ;
String text = result . getString ( " text " ) ;
JSONObject analysisInfo = result . getJSONObject ( " analysis_info " ) ;
String role = analysisInfo . getString ( " role " ) ;
String title = " " ;
if ( role . equals ( " AGENT " ) ) {
title = " 顾问 " ;
} else {
title = " 客户 " ;
}
// 拼接 role 和 text
String chat = title + " : " + text ;
runMaskingRuleInput . setOldStr ( chat ) ;
String corpusChat = dataMaskingRuleService . runMaskingRule ( runMaskingRuleInput ) ;
chatList . append ( corpusChat ) . append ( " \ n " ) ;
inputMap . put ( " chat " , chatList . toString ( ) ) ;
inputMap . put ( " model " , getCarModelList ( ) ) ;
inputMap . put ( " recordId " , aicorpusTelephone . getSourceId ( ) ) ;
inputMap . put ( " version " , 2 ) ;
} ) ;
diFyImageReq . setInputs ( inputMap ) ;
CorpusReportDTO corpusReportDTO = new CorpusReportDTO ( ) ;
corpusReportDTO . setCorpusTime ( formattedDateStartTime ) ;
corpusReportDTO . setRecordId ( aicorpusTelephone . getSourceId ( ) ) ;
corpusReportDTO . setAnalysisScene ( 2l ) ;
ZonedDateTime zonedDateTime = ZonedDateTime . parse ( jsonObject . getString ( " start_time " ) ) ;
DateTimeFormatter formatter = DateTimeFormatter . ofPattern ( " yyyy-MM-dd HH:mm:ss " ) ;
String formattedDateStartTime = zonedDateTime . format ( formatter ) ;
long startTime = System . currentTimeMillis ( ) ;
CompletableFuture . runAsync ( ( ) - > {
try {
log . info ( " 铭牌语料可用许可授权数,总结和分类场景={} " , semaphore . availablePermits ( ) ) ;
// 获取许可 - 如果没有可用许可会阻塞等待
semaphore . acquire ( ) ;
diFyImageReq . setFlowId ( telephoneToken ) ;
log . info ( " 铭牌语料telephoneToken: {} " , telephoneToken ) ;
JSONObject execDifyFlow = diFyService . executeDifyFlow ( diFyImageReq , BusinessTypeEnum . SMART_ASSISTANT . getCode ( ) ,
JSONObject . toJSONString ( corpusReportDTO ) , aicorpusTelephone . getAiAnalysisRequestId ( ) ) ;
inputMap . put ( " chat " , chatList . toString ( ) ) ;
inputMap . put ( " model " , getCarModelList ( ) ) ;
diFyImageReq . s etInputs ( inputMap ) ;
CorpusReportDTO corpusReportDTO = new CorpusReportDTO ( ) ;
corpusReportDTO . setCorpusTime ( formattedDateStartTim e ) ;
corpusReportDTO . setRecordId ( aicorpusTelephone . getSourceId ( ) ) ;
corpusReportDTO . setAnalysisScene ( 2l ) ;
log . info ( " dcc总结场景,总结和分类 runDify execDifyFlow 返回 : {} " , execDifyFlow ) ;
parseDfiyResult ( execDifyFlow , aicorpusTelephone . getSourceId ( ) , formattedDateStartTime , aicorpusTelephone . getAiAnalysisRequestId ( ) ,
BusinessTypeEnum . SMART_ASSISTANT . g etCode ( ) ) ;
} catch ( Exception e ) {
log . error ( " runDify 异常 " , e ) ;
} finally {
// 释放许可
semaphore . release ( ) ;
}
} , executor ) ;
long endTime = System . currentTimeMillis ( ) ;
log . info ( " 第一个业务场景( dcc总结和分类) 执行时间: {} ms " , ( endTime - startTime ) ) ;
//第一个业务场景, 结束
long startTime = System . currentTimeMillis ( ) ;
CompletableFuture . runAsync ( ( ) - > {
try {
log . info ( " 铭牌语料可用许可授权数,总结和分类 场景={} " , semaphore . availablePermits ( ) ) ;
// 获取许可 - 如果没有可用许可会阻塞等待
semaphore . acquire ( ) ;
diFyImageReq . setFlowId ( telephoneToken ) ;
log . info ( " 铭牌语料telephoneToken: {} " , telephoneToken ) ;
JSONObject execDifyFlow = diFyService . executeDifyFlow ( diFyImageReq , BusinessTypeEnum . SMART_ASSISTANT . getCode ( ) ,
JSONObject . toJSONString ( corpusReportDTO ) , aicorpusTelephone . getAiAnalysisRequestId ( ) ) ;
long startTime2 = System . currentTimeMillis ( ) ;
CompletableFuture . runAsync ( ( ) - > {
try {
log . info ( " 铭牌语料可用许可授权数,客户画像 场景={} " , semaphore . availablePermits ( ) ) ;
Thread . sleep ( 2000 ) ;
// 获取许可 - 如果没有可用许可会阻塞等待
semaphore . acquire ( ) ;
// 创建新的DiFyReq对象以避免线程安全问题
inputMap . put ( " businessId " , aicorpusTelephone . getSourceId ( ) ) ;
inputMap . put ( communicateDateStr , aicorpusTelephone . getTranscribeTime ( ) ) ;
inputMap . put ( " analysisScene " , " 2 " ) ;
diFyImageReq . setInputs ( inputMap ) ;
log . info ( " dcc总结场景,总结和分类 runDify execDifyFlow 返回 : {} " , execDifyFlow ) ;
parseDfiyResult ( execDifyFlow , aicorpusTelephone . getSourceId ( ) , formattedDateStartTime , aicorpusTelephone . getAiAnalysisRequestId ( ) ,
BusinessTypeEnum . SMART_ASSISTANT . getCode ( ) ) ;
} catch ( Exception e ) {
log . error ( " runDify 异常 " , e ) ;
} finally {
// 释放许可
semaphore . release ( ) ;
}
} , executor ) ;
long endTime = System . currentTimeMillis ( ) ;
log . info ( " 第一个业务场景( dcc总结和分类) 执行时间: {} ms " , ( endTime - startTime ) ) ;
//第一个业务场景, 结束
diFyImageReq . setFlowId ( oneTokenPortrait ) ;
log . info ( " dcc总结场景,客户画像runDify execDifyFlow 返回 : {} " , oneTokenPortrait ) ;
JSONObject execDifyFlowForPortrait = diFyService . executeDifyFlow ( diFyImageReq , BusinessTypeEnum . CORPUS_PORTRAIT_DCC . getCode ( ) ,
JSONObject . toJSONString ( corpusReportDTO ) , aicorpusTelephone . getAiAnalysisRequestId ( ) ) ;
long startTime2 = System . currentTimeMillis ( ) ;
CompletableFuture . runAsync ( ( ) - > {
try {
log . info ( " 铭牌语料可用许可授权数,总结和分类场景={} " , semaphore . availablePermits ( ) ) ;
// 获取许可 - 如果没有可用许可会阻塞等待
semaphore . acquire ( ) ;
// 创建新的DiFyReq对象以避免线程安全问题
inputMap . put ( " businessId " , aicorpusTelephone . getSourceId ( ) ) ;
inputMap . put ( " communicateDate " , aicorpusTelephone . getTranscribeTime ( ) ) ;
inputMap . put ( " analysisScene " , " 2 " ) ;
diFyImageReq . setInputs ( inputMap ) ;
log . info ( " dcc客户画像场景 runDify execDifyFlow 返回 , dcc: {} " , execDifyFlowForPortrait ) ;
diFyImageReq . s etFlowId ( oneTokenPortrait ) ;
log . info ( " dcc总结场景,客户画像runDify execDifyFlow 返回 : {} " , oneTokenPortrait ) ;
JSONObject execDifyFlowForPortrait = diFyService . executeDifyFlow ( diFyImageReq , BusinessTypeEnum . CORPUS_PORTRAIT_DCC . getCode ( ) ,
JSONObject . toJSONString ( corpusReportDTO ) , aicorpusTelephone . getAiAnalysisRequestId ( ) ) ;
log . info ( " dcc客户画像场景 runDify execDifyFlow 返回 , dcc: {} " , execDifyFlowForPortrait ) ;
parseDfiyResult ( execDifyFlowForPortrait , aicorpusTelephone . getSourceId ( ) , formattedDateStartTime ,
aicorpusTelephone . getAiAnalysisRequestId ( ) , BusinessTypeEnum . CORPUS_PORTRAIT_DCC . getCode ( ) ) ;
} catch ( Exception e ) {
log . error ( " runDify 异常 " , e ) ;
} finally {
// 释放许可
semaphore . release ( ) ;
}
} , executor ) ;
long endTime2 = System . currentTimeMillis ( ) ;
log . info ( " 第二个业务场景( dcc用户画像) 执行时间: {} ms " , ( endTime2 - startTime2 ) ) ;
//第二个业务场景, 结束
parseDfiyResult ( execDifyFlowForPortrait , aicorpusTelephone . g etSourceId ( ) , formattedDateStartTime ,
aicorpusTelephone . getAiAnalysisRequestId ( ) , BusinessTypeEnum . CORPUS_PORTRAIT_DCC . getCode ( ) ) ;
} catch ( Exception e ) {
log . error ( " runDify 异常 " , e ) ;
} finally {
// 释放许可
semaphore . release ( ) ;
}
} , executor ) ;
long endTime2 = System . currentTimeMillis ( ) ;
log . info ( " 第二个业务场景( dcc用户画像) 执行时间: {} ms " , ( endTime2 - startTime2 ) ) ;
//第二个业务场景, 结束
}
} catch ( Exception e ) {
log . error ( " runTelephoneCorpusDify error:{} " , e ) ;
}
}
private void parseDfiyResult ( JSONObject execDifyFlow , String recordId , String communicateDate ,
String aiAnalysisRequestIdDB , String businessType ) {
log . info ( " dcc语料 parseDfiyResult , businessType: {} execDifyFlow: {} " , businessType , execDifyFlow ) ;
log . info ( " dcc语料 parseDfiyResult , businessType: {} execDifyFlow: {}, recordId: {},communicateDate: {} " , businessType , execDifyFlow , recordId , communicateDate );
try {
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 ( " 电话语料解析为空, text:{} " , text ) ;
return ;
}
String aiAnalysisRequestId = execDifyFlow . getString ( " aiAnalysisRequestId " ) ;
Map < String , String > ltoMap = new HashMap ( ) ;
ltoMap . put ( " analysisRecordId " , aiAnalysisRequestId ) ;
ltoMap . put ( " analysisScene " , " 2 " ) ;
ltoMap . put ( " recordId " , recordId ) ;
ltoMap . put ( " communicateDate " , communicateDate ) ;
ltoMap . put ( " analysisResult " , resultStrOne ) ;
ltoMap . put ( " analysisDetail " , resultStrTwo ) ;
JSONObject text = execDifyFlow . getJSONObject ( outputsStr ) ;
// 发送MQ
if ( BusinessTypeEnum . CORPUS_PORTRAIT_DCC . getCode ( ) . equals ( businessType ) ) { //DCC 客户画像
log . info ( " send mq 电话语料场景,客户画像需求, {} " , ltoMap ) ;
sendMq ( CategoryEnum . PORTRAIT_ALLIN . getCode ( ) , JSONObjec t. toJSONString ( ltoMap ) ) ;
log . info ( " send mq 电话语料场景,客户画像需求, {} " , text ) ;
sendMq ( CategoryEnum . PORTRAIT_ALLIN . getCode ( ) , tex t. toJSONString ( ) ) ;
} else { // dcc 总结
log . info ( " send mq 电话语料场景,总结需求, {} " , ltoMap ) ;
sendMq ( CategoryEnum . PHONE_VOICE . getCode ( ) , JSONObjec t. toJSONString ( ltoMap ) ) ;
log . info ( " send mq 电话语料场景,总结需求, {} " , text ) ;
sendMq ( CategoryEnum . PHONE_VOICE . getCode ( ) , tex t. toJSONString ( ) ) ;
}
try {
if ( StringUtils . isNotEmpty ( aiAnalysisRequestIdDB ) ) {
@@ -260,10 +254,10 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
aiAnalysisErrors . setAiAnalysisErrorHandlingStatus ( " 1 " ) ;
aiAnalysisErrorsService . updateAiAnalysisErrors ( aiAnalysisErrors ) ;
}
aiAnalysisRequestLogsService . saveOrUpdateAiAnalysisRequestLogs ( AiAnalysisRequestLogs . builder ( ) . aiAnalysisRequestId ( execDifyFlow . getString ( " aiAnalysisRequestId" ) ) . businessResponse ( JSONObjec t. toJSONString ( ltoMap ) ) . build ( ) ) ;
aiAnalysisRequestLogsService . saveOrUpdateAiAnalysisRequestLogs ( AiAnalysisRequestLogs . builder ( ) . aiAnalysisRequestId ( execDifyFlow . getString ( aiAnalysisRequestIdStr ) ) . businessResponse ( tex t. toJSONString ( ) ) . build ( ) ) ;
} catch ( Exception e ) {
log . info ( " 电话语料处理保存报告异常processItem: {} " , e ) ;
log . info ( " 电话语料处理保存报告异常processItem: {} " , e . getMessage ( ) );
}
} else {
log . info ( " dcc语料 parseDfiyResult 非正常状态 " ) ;
@@ -346,7 +340,7 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
String newAiAnalysisRequestId = StringUtils . isEmpty ( oldAiAnalysisRequestId ) ? AiAnalysisUtils . getAiAnalysisRequestId ( BusinessTypeEnum . SMART_ASSISTANT_4IN1 . getCode ( ) ) :
oldAiAnalysisRequestId ;
inputMap . put ( " aiAnalysisRequestId" , newAiAnalysisRequestId ) ;
inputMap . put ( aiAnalysisRequestIdStr , newAiAnalysisRequestId ) ;
diFyImageReq . setInputs ( inputMap ) ;
CorpusReportDTO corpusReportDTO = new CorpusReportDTO ( ) ;
corpusReportDTO . setRecordId ( aicorpusTelephone . getSourceId ( ) ) ;
@@ -356,9 +350,9 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
// 获取配置
JSONObject execDifyFlow = diFyService . executeDifyFlow ( diFyImageReq , BusinessTypeEnum . SMART_ASSISTANT_4IN1 . getCode ( ) , JSONObject . toJSONString ( corpusReportDTO ) , newAiAnalysisRequestId ) ;
log . info ( " runDify智能客服4IN1 {} " , execDifyFlow ) ;
String aiAnalysisRequestId = execDifyFlow . getString ( " aiAnalysisRequestId" ) ;
String aiAnalysisRequestId = execDifyFlow . getString ( aiAnalysisRequestIdStr ) ;
if ( null ! = execDifyFlow & & execDifyFlow . get ( " status " ) . equals ( " succeeded " ) ) {
JSONObject outputsJson = execDifyFlow . getJSONObject ( " outputs" ) ;
JSONObject outputsJson = execDifyFlow . getJSONObject ( outputsStr ) ;
String message = JSONObject . toJSONString ( outputsJson ) ;
log . info ( " send 4IN1mq {} " , message ) ;
intelligentCustomerService . sendMq ( message ) ;
@@ -368,7 +362,7 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
aiAnalysisErrors . setAiAnalysisErrorHandlingStatus ( " 1 " ) ;
aiAnalysisErrorsService . updateAiAnalysisErrors ( aiAnalysisErrors ) ;
}
aiAnalysisRequestLogsService . saveOrUpdateAiAnalysisRequestLogs ( AiAnalysisRequestLogs . builder ( ) . aiAnalysisRequestId ( execDifyFlow . getString ( " aiAnalysisRequestId" ) ) . businessResponse ( message ) . build ( ) ) ;
aiAnalysisRequestLogsService . saveOrUpdateAiAnalysisRequestLogs ( AiAnalysisRequestLogs . builder ( ) . aiAnalysisRequestId ( execDifyFlow . getString ( aiAnalysisRequestIdStr ) ) . businessResponse ( message ) . build ( ) ) ;
}
@@ -402,7 +396,7 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
@Override
public void sendMq ( String tag , String message ) {
log . info ( " RocketMQ主题和标签: " + topic + " : " + tag ) ;
log . info ( " RocketMQ主题和标签: " + topic + " : " + tag + message ) ;
rocketMqTemplate . asyncSend ( topic + " : " + tag , MessageBuilder . withPayload ( message ) . build ( ) ,
new SendCallback ( ) {
@Override