Merge branch 'dev-feature-20250923-json' into 'feature-20250923'
Dev feature 20250923 json See merge request nsc/chinavolvoaidrives/ai-analytic-center!55
This commit is contained in:
@@ -15,6 +15,7 @@ import org.springframework.kafka.support.KafkaHeaders;
|
|||||||
import org.springframework.messaging.handler.annotation.Header;
|
import org.springframework.messaging.handler.annotation.Header;
|
||||||
import org.springframework.stereotype.Component;
|
import org.springframework.stereotype.Component;
|
||||||
import org.springframework.web.bind.annotation.RestController;
|
import org.springframework.web.bind.annotation.RestController;
|
||||||
|
import org.springframework.beans.factory.annotation.Value;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* @ClassName NameplateKafkaConsumer
|
* @ClassName NameplateKafkaConsumer
|
||||||
@@ -33,16 +34,20 @@ public class NameplateKafkaConsumer {
|
|||||||
@Autowired
|
@Autowired
|
||||||
private TmNameplateCorpusService tmNameplateCorpusService;
|
private TmNameplateCorpusService tmNameplateCorpusService;
|
||||||
|
|
||||||
|
@Value("${analyticCenterKafka.consumer.concurrency}")
|
||||||
|
private String concurrency;
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
@KafkaListener(topics = "${analyticCenterKafka.consumer.topic}",
|
@KafkaListener(topics = "${analyticCenterKafka.consumer.topic}",
|
||||||
groupId = "${analyticCenterKafka.consumer.group}" ,
|
groupId = "${analyticCenterKafka.consumer.group}" ,
|
||||||
containerFactory = "analyticCenterConsumerFactory",
|
containerFactory = "analyticCenterConsumerFactory",
|
||||||
concurrency = "3")
|
concurrency = "${analyticCenterKafka.consumer.concurrency}")
|
||||||
public void listen(String recordMessages, Acknowledgment ack,
|
public void listen(String recordMessages, Acknowledgment ack,
|
||||||
@Header(KafkaHeaders.RECEIVED_PARTITION_ID) Integer partitionId,
|
@Header(KafkaHeaders.RECEIVED_PARTITION_ID) Integer partitionId,
|
||||||
@Header(KafkaHeaders.OFFSET) Long offset) {
|
@Header(KafkaHeaders.OFFSET) Long offset) {
|
||||||
long startTime = System.currentTimeMillis();
|
long startTime = System.currentTimeMillis();
|
||||||
log.info("nameplateKafkaConsumer 当前线程: {}, 线程ID: {},计数:{}", Thread.currentThread().getName(), Thread.currentThread().getId());
|
log.info("nameplateKafkaConsunameplateKafkaConsumermer 当前线程: {}, 线程ID: {},计数:{}", Thread.currentThread().getName(), Thread.currentThread().getId(),concurrency);
|
||||||
log.info("nameplateKafkaConsumerMessage: {}", recordMessages);
|
log.info("nameplateKafkaConsumerMessage: {}", recordMessages);
|
||||||
// 初始化绑定 Consumer
|
// 初始化绑定 Consumer
|
||||||
if (StringUtils.isEmpty(recordMessages)) {
|
if (StringUtils.isEmpty(recordMessages)) {
|
||||||
|
|||||||
@@ -49,7 +49,7 @@ public class DataMaskingRuleServiceImpl extends ServiceImpl<DataMaskingRuleMappe
|
|||||||
log.info("runMaskingRule categoryName:{} not support", ruleItem.getCategoryName());
|
log.info("runMaskingRule categoryName:{} not support", ruleItem.getCategoryName());
|
||||||
}
|
}
|
||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
e.printStackTrace();
|
log.error("runMaskingRule errorMessage"+e.getMessage());
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user