修改kafka参数

This commit is contained in:
zren25
2025-05-16 14:44:05 +08:00
parent f7bd2c2708
commit 33d22a3e31

View File

@@ -19,7 +19,9 @@ import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.RestController; import org.springframework.web.bind.annotation.RestController;
import javax.annotation.Resource; import javax.annotation.Resource;
import java.util.concurrent.*; import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
/** /**
* @ClassName NameplateKafkaConsumer * @ClassName NameplateKafkaConsumer
@@ -49,10 +51,11 @@ public class NameplateKafkaConsumer {
private int messageCount = 0; private int messageCount = 0;
@KafkaListener(topics = "${analyticCenterKafka.consumer.topic}", @KafkaListener(topics = "${analyticCenterKafka.consumer.topic}",
groupId = "${analyticCenterKafka.consumer.group}" , groupId = "${analyticCenterKafka.consumer.group}" ,
containerFactory = "analyticCenterConsumerFactory") containerFactory = "analyticCenterConsumerFactory",
concurrency = "3")
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) throws InterruptedException { @Header(KafkaHeaders.OFFSET) Long offset) {
long startTime = System.currentTimeMillis(); long startTime = System.currentTimeMillis();
messageCount++; messageCount++;
@@ -67,7 +70,6 @@ public class NameplateKafkaConsumer {
NameplateTableKafkaDTO tmNameplateCorpus = JSON.parseObject(recordMessages, NameplateTableKafkaDTO.class); NameplateTableKafkaDTO tmNameplateCorpus = JSON.parseObject(recordMessages, NameplateTableKafkaDTO.class);
log.info("nameplateKafkaConsumerParseType: {},size:{}", tmNameplateCorpus.getType(), tmNameplateCorpus.getData().size()); log.info("nameplateKafkaConsumerParseType: {},size:{}", tmNameplateCorpus.getType(), tmNameplateCorpus.getData().size());
if(tmNameplateCorpus.getType().equals("INSERT")){ if(tmNameplateCorpus.getType().equals("INSERT")){
CountDownLatch latch = new CountDownLatch(tmNameplateCorpus.getData().size());
tmNameplateCorpus.getData().forEach(nameplate -> tmNameplateCorpus.getData().forEach(nameplate ->
CompletableFuture.runAsync(() -> { CompletableFuture.runAsync(() -> {
try { try {
@@ -76,18 +78,14 @@ public class NameplateKafkaConsumer {
} catch (Exception e) { } catch (Exception e) {
log.error("nameplateKafkaConsumer铭牌解析处理出错 {}:{}ex:{}", log.error("nameplateKafkaConsumer铭牌解析处理出错 {}:{}ex:{}",
partitionId, offset, e.getMessage()); partitionId, offset, e.getMessage());
}finally {
log.info("nameplateKafkaConsumerack.acknowledge");
latch.countDown();
} }
}, executor)); }, executor));
try { try {
log.info("nameplateKafkaConsumerack.acknowledge:{}");
ack.acknowledge(); // 确保无论成功与否都提交 offset ack.acknowledge(); // 确保无论成功与否都提交 offset
} catch (IllegalStateException e) { } catch (IllegalStateException e) {
log.warn("Offset 已提交,跳过重复提交"); log.warn("Offset 已提交,跳过重复提交");
} }
boolean allDone = latch.await(2, TimeUnit.MINUTES);
log.info("nameplateKafkaConsumerack.acknowledge:{}",allDone);
} }
}else { }else {
ack.acknowledge(); // 空消息直接跳过 ack.acknowledge(); // 空消息直接跳过