增加日志
This commit is contained in:
@@ -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<ConsumerRecord<String, Object>> recordMessages) {
|
||||
public void listen( @Payload List<String> messageList,
|
||||
@Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) List<String> keys,
|
||||
@Header(KafkaHeaders.RECEIVED_PARTITION_ID) List<Integer> partitions,
|
||||
@Header(KafkaHeaders.OFFSET) List<Long> 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<CompletableFuture<Void>> 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<String, Object> 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<String, Object> 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);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user