增加耗时时间日志

This commit is contained in:
zren25
2025-03-31 14:19:50 +08:00
parent 1ce61d4bfc
commit 2761de1227

View File

@@ -9,7 +9,6 @@ import com.volvo.ai.analytic.center.entity.TmTelephoneCorpus;
import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService; import com.volvo.ai.analytic.center.service.TmTelephoneCorpusService;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.apache.commons.collections.CollectionUtils; import org.apache.commons.collections.CollectionUtils;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.beans.BeanUtils; import org.springframework.beans.BeanUtils;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.annotation.KafkaListener;
@@ -44,6 +43,7 @@ public class CorpusProcessKafkaProducer {
@PostMapping("corpusProcessKafkaConsumer") @PostMapping("corpusProcessKafkaConsumer")
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}") @KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}")
public void listen(List<String> recordMessages) { public void listen(List<String> recordMessages) {
long startTime = System.currentTimeMillis();
try { try {
log.info("CorpusProcessKafkaProducer Received message: {}", recordMessages); log.info("CorpusProcessKafkaProducer Received message: {}", recordMessages);
// 获取消息列表 // 获取消息列表
@@ -65,18 +65,18 @@ public class CorpusProcessKafkaProducer {
tmTelephoneCorpus.setCreateBy("kafka"); tmTelephoneCorpus.setCreateBy("kafka");
tmTelephoneCorpus.setCreateTime(LocalDateTime.now()); tmTelephoneCorpus.setCreateTime(LocalDateTime.now());
ExecutorService runTelephoneExecutor = Executors.newFixedThreadPool(optimalThreadPoolSize); ExecutorService runTelephoneExecutor = Executors.newFixedThreadPool(optimalThreadPoolSize);
// 条件只处理dcc的 10s通话时间以上 // 条件只处理dcc的 10s通话时间以上
log.info("CorpusProcessKafkaProducer getCategoryCode: {}audioFileId{}sourceId{}", tmTelephoneCorpus.getCategoryCode(), aicorpusTelephone.getAudioFileId(), aicorpusTelephone.getSourceId()); log.info("CorpusProcessKafkaProducer getCategoryCode: {}audioFileId{}sourceId{}", tmTelephoneCorpus.getCategoryCode(), aicorpusTelephone.getAudioFileId(), aicorpusTelephone.getSourceId());
if (Constant.CHANNEL_DCC.equals(tmTelephoneCorpus.getCategoryCode())) { if (Constant.CHANNEL_DCC.equals(tmTelephoneCorpus.getCategoryCode())) {
log.info(" dcc 语料开始处理: {}"); log.info(" dcc 语料开始处理");
tmTelephoneCorpusService.saveTelephoneCorpus(tmTelephoneCorpus); tmTelephoneCorpusService.saveTelephoneCorpus(tmTelephoneCorpus);
long startTimeDify = System.currentTimeMillis();
CompletableFuture.runAsync(() -> { CompletableFuture.runAsync(() -> {
tmTelephoneCorpusService.runTelephoneCorpusDify(aicorpusTelephone); tmTelephoneCorpusService.runTelephoneCorpusDify(aicorpusTelephone);
}, runTelephoneExecutor); }, runTelephoneExecutor);
// 关闭线程池 // 关闭线程池
runTelephoneExecutor.shutdown(); runTelephoneExecutor.shutdown();
log.info(" dify处理耗时{}", System.currentTimeMillis() - startTimeDify);
} }
} catch (Exception e) { } catch (Exception e) {
log.error("CorpusProcessKafkaProducer 电话语料 解析JSON出错: {}", e.getMessage()); log.error("CorpusProcessKafkaProducer 电话语料 解析JSON出错: {}", e.getMessage());
@@ -87,7 +87,7 @@ public class CorpusProcessKafkaProducer {
executor.shutdown(); executor.shutdown();
} }
log.info("Kafka 消息处理完成,耗时:{}", System.currentTimeMillis() - startTime);
// 在这里可以添加对解析后的对象的进一步处理逻辑 // 在这里可以添加对解析后的对象的进一步处理逻辑
} catch (Exception e) { } catch (Exception e) {
log.error("CorpusProcessKafkaProducer 电话语料 解析JSON出错: {}" , e.getMessage()); log.error("CorpusProcessKafkaProducer 电话语料 解析JSON出错: {}" , e.getMessage());