diff --git a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/TestController.java b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/TestController.java index 32389a6..8a1f934 100644 --- a/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/TestController.java +++ b/ai-analytic-center-biz/src/main/java/com/volvo/ai/analytic/center/controller/TestController.java @@ -73,7 +73,9 @@ public class TestController { public ResultMsg mockDccKafka(@RequestBody AicorpusTelephoneDTO message) { for(int i=0;i<20;i++){ - message.setSourceId("000001"+i); + String sourceId = "2000000"+i; + message.setSourceId(sourceId); + log.info("发送消息:"+sourceId); dccKafkaProducer.send("topic_voc_covert_text_log",JSONObject.toJSONString(message)); } return ResultMsg.ok("ok"); 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 213d3ab..990aa9d 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 @@ -17,8 +17,6 @@ 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.kafka.annotation.KafkaListener; -import org.springframework.kafka.support.KafkaHeaders; -import org.springframework.messaging.handler.annotation.Header; import org.springframework.messaging.handler.annotation.Payload; import org.springframework.messaging.support.MessageBuilder; import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; @@ -61,10 +59,7 @@ public class CorpusProcessKafkaProducer { @KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}") - public void listen(@Payload List messageList, - @Header(KafkaHeaders.RECEIVED_PARTITION_ID) List partitions, - @Header(KafkaHeaders.OFFSET) List offsets - ) { + public void listen(@Payload List messageList) { long startTime = System.currentTimeMillis(); log.info("corpusProcessKafkaProducerMessageSize:{}", messageList.size()); if (CollectionUtils.isEmpty(messageList)) { @@ -72,13 +67,8 @@ public class CorpusProcessKafkaProducer { return; } try { - for (int i = 0; i < messageList.size(); i++) { - String message = messageList.get(i); - log.debug("corpusProcessKafkaProducer message:{}", message); - int partition = partitions.get(i); - long offset = offsets.get(i); - - log.debug("corpusProcessKafkaProducer message: topic={}, partition={}, offset={}", "your-topic",partition, offset); + for (String message:messageList) { + log.debug("corpusProcessKafkaProducerMessage:{}", message); // 提交异步任务 CompletableFuture.runAsync(() -> { @@ -125,8 +115,7 @@ public class CorpusProcessKafkaProducer { * 处理 DCC 语料逻辑(异步发送 MQ) */ private void handleDccCorpus(TmTelephoneCorpus corpus, String rawMessage) { - log.info("handleDccCorpus: categoryCode={}, sourceId={}", - corpus.getCategoryCode(), corpus.getSourceId()); + log.info("corpusProcessKafkaProducerSaveSourceId: {}",corpus.getSourceId()); // 4. 保存语料 tmTelephoneCorpusService.saveTelephoneCorpus(corpus);