修改kafka监听对象
This commit is contained in:
@@ -43,22 +43,19 @@ public class CorpusProcessKafkaProducer {
|
||||
|
||||
@PostMapping("corpusProcessKafkaConsumer")
|
||||
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}")
|
||||
public void listen(List<ConsumerRecord<String, Object>> recordMessage) {
|
||||
public void listen(List<String> recordMessages) {
|
||||
try {
|
||||
log.info("CorpusProcessKafkaProducer Received message: {}", recordMessage);
|
||||
log.info("CorpusProcessKafkaProducer Received message: {}", recordMessages);
|
||||
// 获取消息列表
|
||||
int optimalThreadPoolSize = Runtime.getRuntime().availableProcessors() + 2;
|
||||
log.info("获取的线程数:{}", optimalThreadPoolSize);
|
||||
ExecutorService executor = Executors.newFixedThreadPool(optimalThreadPoolSize);
|
||||
if(CollectionUtils.isNotEmpty(recordMessage)){
|
||||
log.info("CorpusProcessKafkaProducer List size: {}", recordMessage.size());
|
||||
for (ConsumerRecord<String, Object> record : recordMessage) {
|
||||
log.info("CorpusProcessKafkaProducer List record.value: {}",record.value());
|
||||
log.info("CorpusProcessKafkaProducer List record.key: {}",record.key());
|
||||
|
||||
if(CollectionUtils.isNotEmpty(recordMessages)){
|
||||
log.info("CorpusProcessKafkaProducer List size: {}", recordMessages.size());
|
||||
for (String message : recordMessages) {
|
||||
log.info("CorpusProcessKafkaProducer message: {}",message);
|
||||
executor.submit(() -> {
|
||||
try {
|
||||
String message = (String) record.value();
|
||||
AicorpusTelephoneDTO aicorpusTelephone = objectMapper.readValue(message, AicorpusTelephoneDTO.class);
|
||||
log.info("aicorpusTelephone categoryCode:{}, display: {}", aicorpusTelephone.getCategoryCode(), aicorpusTelephone.getDisplay());
|
||||
DisplayDTO display = objectMapper.readValue(aicorpusTelephone.getDisplay(), DisplayDTO.class);
|
||||
|
||||
Reference in New Issue
Block a user