加日志 ,修改并发
This commit is contained in:
@@ -34,6 +34,7 @@ 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.messaging.support.MessageBuilder;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
import org.springframework.stereotype.Service;
|
||||
import org.springframework.transaction.annotation.Transactional;
|
||||
|
||||
@@ -46,6 +47,7 @@ import java.util.Map;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.Semaphore;
|
||||
import java.util.stream.Collectors;
|
||||
|
||||
|
||||
@@ -95,6 +97,13 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
|
||||
@Autowired
|
||||
CorpusQuestionProducer corpusQuestionProducer;
|
||||
|
||||
// 使用多线程并行执行两个业务场景
|
||||
@Autowired
|
||||
@Resource(name = "threadPoolTaskExecutor")
|
||||
private ThreadPoolTaskExecutor executor;
|
||||
|
||||
private Semaphore semaphore = new Semaphore(6); // 限制并发数
|
||||
|
||||
@Override
|
||||
@Transactional
|
||||
public void saveTelephoneCorpus(TmTelephoneCorpus tmTelephoneCorpus) {
|
||||
@@ -157,12 +166,12 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
|
||||
corpusReportDTO.setRecordId(aicorpusTelephone.getSourceId());
|
||||
corpusReportDTO.setAnalysisScene(2l);
|
||||
|
||||
// 使用多线程并行执行两个业务场景
|
||||
ExecutorService executorService = Executors.newFixedThreadPool(2);
|
||||
|
||||
long startTime = System.currentTimeMillis();
|
||||
CompletableFuture<Void> summaryFuture = CompletableFuture.runAsync(() -> {
|
||||
CompletableFuture.runAsync(() -> {
|
||||
try {
|
||||
log.info("铭牌语料可用许可授权数,总结和分类场景={}", semaphore.availablePermits());
|
||||
// 获取许可 - 如果没有可用许可会阻塞等待
|
||||
semaphore.acquire();
|
||||
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyImageReq, BusinessTypeEnum.SMART_ASSISTANT.getCode(),
|
||||
JSONObject.toJSONString(corpusReportDTO), aicorpusTelephone.getAiAnalysisRequestId());
|
||||
|
||||
@@ -171,15 +180,21 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
|
||||
BusinessTypeEnum.SMART_ASSISTANT.getCode());
|
||||
} catch (Exception e) {
|
||||
log.error("runDify 异常", e);
|
||||
}finally {
|
||||
// 释放许可
|
||||
semaphore.release();
|
||||
}
|
||||
}, executorService);
|
||||
}, executor);
|
||||
long endTime = System.currentTimeMillis();
|
||||
log.info("第一个业务场景(dcc总结和分类)执行时间: {} ms", (endTime - startTime));
|
||||
//第一个业务场景, 结束
|
||||
|
||||
long startTime2 = System.currentTimeMillis();
|
||||
CompletableFuture<Void> portraitFuture = CompletableFuture.runAsync(() -> {
|
||||
CompletableFuture.runAsync(() -> {
|
||||
try {
|
||||
log.info("铭牌语料可用许可授权数,总结和分类场景={}", semaphore.availablePermits());
|
||||
// 获取许可 - 如果没有可用许可会阻塞等待
|
||||
semaphore.acquire();
|
||||
// 创建新的DiFyReq对象以避免线程安全问题
|
||||
inputMap.put("businessId",aicorpusTelephone.getSourceId());
|
||||
inputMap.put("communicateDate",aicorpusTelephone.getTranscribeTime());
|
||||
@@ -196,8 +211,11 @@ public class TmTelephoneCorpusServiceImpl extends ServiceImpl<TmTelephoneCorpusM
|
||||
aicorpusTelephone.getAiAnalysisRequestId(), BusinessTypeEnum.CORPUS_PORTRAIT_DCC.getCode());
|
||||
} catch (Exception e) {
|
||||
log.error("runDify 异常", e);
|
||||
}finally {
|
||||
// 释放许可
|
||||
semaphore.release();
|
||||
}
|
||||
}, executorService);
|
||||
}, executor);
|
||||
long endTime2 = System.currentTimeMillis();
|
||||
log.info("第二个业务场景(dcc用户画像)执行时间: {} ms", (endTime2 - startTime2));
|
||||
//第二个业务场景, 结束
|
||||
|
||||
Reference in New Issue
Block a user