From 8bc79adc359394f855180bdb53c7067d30ba4342 Mon Sep 17 00:00:00 2001 From: zren25 Date: Tue, 27 May 2025 17:04:16 +0800 Subject: [PATCH] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=E6=97=A5=E5=BF=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../center/mq/CorpusProcessKafkaProducer.java | 48 ++++++++++--------- 1 file changed, 26 insertions(+), 22 deletions(-) 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 433656e..571caf7 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 @@ -5,12 +5,10 @@ import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import com.volvo.ai.analytic.center.constant.Constant; import com.volvo.ai.analytic.center.dto.corpus.AicorpusTelephoneDTO; -import com.volvo.ai.analytic.center.dto.corpus.DisplayDTO; 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.apache.rocketmq.client.producer.SendCallback; import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.spring.core.RocketMQTemplate; @@ -19,6 +17,9 @@ 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; import org.springframework.stereotype.Component; @@ -26,7 +27,6 @@ import org.springframework.web.bind.annotation.RestController; import javax.annotation.Resource; import java.time.LocalDateTime; -import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ThreadPoolExecutor; @@ -62,22 +62,29 @@ public class CorpusProcessKafkaProducer { private ThreadPoolTaskExecutor executor; @KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}") - public void listen(List> recordMessages) { + public void listen( @Payload List messageList, + @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) List keys, + @Header(KafkaHeaders.RECEIVED_PARTITION_ID) List partitions, + @Header(KafkaHeaders.OFFSET) List offsets + ) { long startTime = System.currentTimeMillis(); - if (CollectionUtils.isEmpty(recordMessages)) { + log.info("corpusProcessKafkaProducerMessageSize:{}", messageList.size()); + if (CollectionUtils.isEmpty(messageList)) { log.debug("Received empty message batch."); return; } - - List> futures = new ArrayList<>(); try { - log.info("Received {} messages. Start processing...", recordMessages.size()); + for (int i = 0; i < messageList.size(); i++) { + String message = messageList.get(i); + log.debug("corpusProcessKafkaProducer message:{}", message); + String key = keys.get(i); + int partition = partitions.get(i); + long offset = offsets.get(i); - for (ConsumerRecord record : recordMessages) { - // 使用 CompletableFuture 提交异步任务 - futures.add( - CompletableFuture.runAsync(() -> processSingleRecord(record), executor) - ); + log.debug("corpusProcessKafkaProducer message:key:{}, topic={}, partition={}, offset={}", "your-topic", key,partition, offset); + + // 提交异步任务 + CompletableFuture.runAsync(() -> processSingleRecord(message), executor); } } catch (Exception e) { @@ -102,13 +109,10 @@ public class CorpusProcessKafkaProducer { /** * 单条消息处理逻辑(解耦核心逻辑) */ - private void processSingleRecord(ConsumerRecord record) { + private void processSingleRecord(String record) { try { - log.debug("Processing message: topic={}, partition={}, offset={}", - record.topic(), record.partition(), record.offset()); - - String message = (String) record.value(); - AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(message, AicorpusTelephoneDTO.class); + log.info("processSingleRecord message:{}", record); + AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(record, AicorpusTelephoneDTO.class); TmTelephoneCorpus tmTelephoneCorpus = new TmTelephoneCorpus(); BeanUtils.copyProperties(aicorpusTelephone, tmTelephoneCorpus); tmTelephoneCorpus.setCreateBy("kafka"); @@ -116,12 +120,12 @@ public class CorpusProcessKafkaProducer { // 3. 逻辑拆分:处理 DCC 语料 if (Constant.CHANNEL_DCC.equals(tmTelephoneCorpus.getCategoryCode())) { - handleDccCorpus(tmTelephoneCorpus, message); + handleDccCorpus(tmTelephoneCorpus, record); } } catch (JsonProcessingException e) { - log.error("JSON parsing failed for message: {}", record.value(), e); + log.error("JSON parsing failed for message: {},{}", record, e); } catch (Exception e) { - log.error("Unexpected error processing message: ", e); + log.error("Unexpected error processing message: {}", e); } }