修改message

This commit is contained in:
zren25
2025-03-31 22:46:10 +08:00
parent de204af6b0
commit 3922fb885f

View File

@@ -9,6 +9,7 @@ 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.apache.rocketmq.client.producer.SendCallback; import org.apache.rocketmq.client.producer.SendCallback;
import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate; import org.apache.rocketmq.spring.core.RocketMQTemplate;
@@ -59,7 +60,7 @@ public class CorpusProcessKafkaProducer {
@PostMapping("corpusProcessKafkaProducer") @PostMapping("corpusProcessKafkaProducer")
@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<ConsumerRecord<String, Object>> recordMessages) {
long startTime = System.currentTimeMillis(); long startTime = System.currentTimeMillis();
try { try {
log.info("CorpusProcessKafkaProducer Received message: {}", recordMessages); log.info("CorpusProcessKafkaProducer Received message: {}", recordMessages);
@@ -69,12 +70,12 @@ public class CorpusProcessKafkaProducer {
ExecutorService executor = Executors.newFixedThreadPool(optimalThreadPoolSize); ExecutorService executor = Executors.newFixedThreadPool(optimalThreadPoolSize);
if(CollectionUtils.isNotEmpty(recordMessages)){ if(CollectionUtils.isNotEmpty(recordMessages)){
log.info("CorpusProcessKafkaProducer List size: {}", recordMessages.size()); log.info("CorpusProcessKafkaProducer List size: {}", recordMessages.size());
for (String message : recordMessages) { for (ConsumerRecord<String, Object> record : recordMessages) {
log.info("CorpusProcessKafkaProducer message: {}",message); log.info("CorpusProcessKafkaProducer message: {}",record);
executor.submit(() -> { executor.submit(() -> {
try { try {
String message = (String) record.value();
AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(message, AicorpusTelephoneDTO.class); AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(message, AicorpusTelephoneDTO.class);
String transcribeTimeStr = aicorpusTelephone.getTranscribeTime(); String transcribeTimeStr = aicorpusTelephone.getTranscribeTime();
if (transcribeTimeStr != null) { if (transcribeTimeStr != null) {
DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"); DateTimeFormatter formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
@@ -91,7 +92,6 @@ public class CorpusProcessKafkaProducer {
BeanUtils.copyProperties(aicorpusTelephone, tmTelephoneCorpus); BeanUtils.copyProperties(aicorpusTelephone, tmTelephoneCorpus);
tmTelephoneCorpus.setCreateBy("kafka"); tmTelephoneCorpus.setCreateBy("kafka");
tmTelephoneCorpus.setCreateTime(LocalDateTime.now()); tmTelephoneCorpus.setCreateTime(LocalDateTime.now());
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())) {