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