添加日志,修改并发逻辑
This commit is contained in:
@@ -19,11 +19,14 @@ import org.springframework.web.bind.annotation.RestController;
|
||||
|
||||
import javax.annotation.Resource;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.Semaphore;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
|
||||
/**
|
||||
* @ClassName NameplateKafkaConsumer
|
||||
* @Description 消费tm_nameplate_corpus表binlog的Kafka消息
|
||||
* @Description 消费tm_nameplate_corpus表binlog的Kafka消息
|
||||
* @Author renzhen
|
||||
* @Date 2025-03-04 10:18
|
||||
* @Version 1.0
|
||||
@@ -35,68 +38,72 @@ import java.util.concurrent.Semaphore;
|
||||
@RefreshScope
|
||||
public class NameplateKafkaConsumer {
|
||||
|
||||
@Autowired
|
||||
private TmNameplateCorpusService tmNameplateCorpusService;
|
||||
@Autowired
|
||||
private TmNameplateCorpusService tmNameplateCorpusService;
|
||||
|
||||
@Autowired
|
||||
@Resource(name = "threadPoolTaskExecutor")
|
||||
private ThreadPoolTaskExecutor executor;
|
||||
private Semaphore semaphore = new Semaphore(10); // 限制并发数
|
||||
@Autowired
|
||||
@Resource(name = "threadPoolTaskExecutor")
|
||||
private ThreadPoolTaskExecutor executor;
|
||||
private Semaphore semaphore = new Semaphore(2); // 限制并发数
|
||||
|
||||
@KafkaListener(topics = "${analyticCenterKafka.consumer.topic}", // = smart_assistant_nameplate_topic
|
||||
groupId = "${analyticCenterKafka.consumer.group}" , //smart_assistant_nameplate_topic_group
|
||||
containerFactory = "analyticCenterConsumerFactory",
|
||||
concurrency = "3")
|
||||
public void listen(String recordMessages, Acknowledgment ack,
|
||||
@Header(KafkaHeaders.RECEIVED_PARTITION_ID) Integer partitionId,
|
||||
@Header(KafkaHeaders.OFFSET) Long offset) {
|
||||
long startTime = System.currentTimeMillis();
|
||||
log.info("nameplateKafkaConsumer 当前线程: {}, 线程ID: {},计数:{}", Thread.currentThread().getName(), Thread.currentThread().getId());
|
||||
log.info("nameplateKafkaConsumerMessage,消息:{}", recordMessages);
|
||||
// 初始化绑定 Consumer
|
||||
if (StringUtils.isEmpty(recordMessages)) {
|
||||
ack.acknowledge();
|
||||
return;
|
||||
}
|
||||
@KafkaListener(topics = "${analyticCenterKafka.consumer.topic}", // = smart_assistant_nameplate_topic
|
||||
groupId = "${analyticCenterKafka.consumer.group}", //smart_assistant_nameplate_topic_group
|
||||
containerFactory = "analyticCenterConsumerFactory",
|
||||
concurrency = "3")
|
||||
public void listen(String recordMessages, Acknowledgment ack,
|
||||
@Header(KafkaHeaders.RECEIVED_PARTITION_ID) Integer partitionId,
|
||||
@Header(KafkaHeaders.OFFSET) Long offset) {
|
||||
|
||||
try {
|
||||
NameplateTableKafkaDTO tmNameplateCorpus = JSON.parseObject(recordMessages, NameplateTableKafkaDTO.class);
|
||||
if (!"INSERT".equals(tmNameplateCorpus.getType()) || CollectionUtils.isEmpty(tmNameplateCorpus.getData())) {
|
||||
ack.acknowledge();
|
||||
}
|
||||
if (StringUtils.isEmpty(recordMessages)) {
|
||||
ack.acknowledge();
|
||||
return;
|
||||
}
|
||||
|
||||
tmNameplateCorpus.getData().forEach(nameplate ->
|
||||
CompletableFuture.runAsync(() -> {
|
||||
try {
|
||||
log.info(String.format("尝试获取许可,当前可用许可数: %d",
|
||||
semaphore.availablePermits()));
|
||||
if (semaphore.availablePermits() == 0)
|
||||
log.info("获取许可 失败,❌ 任务被中断");
|
||||
semaphore.acquire(); // 获取许可
|
||||
log.info(String.format("获取许可✅ 成功后,当前可用许可数:: %d",
|
||||
semaphore.availablePermits()));
|
||||
tmNameplateCorpusService.processItem(nameplate);
|
||||
} catch (Exception e) {
|
||||
log.error("corpusPortrait画像铭牌异步任务执行失败", e);
|
||||
} finally {
|
||||
log.info("corpusPortrait 释放锁,availablePermits {}", semaphore.availablePermits());
|
||||
try {
|
||||
ack.acknowledge(); // 提交 offset
|
||||
} catch (IllegalStateException e) {
|
||||
log.warn("Offset 已提交,跳过重复提交");
|
||||
}
|
||||
semaphore.release(); // 释放许可
|
||||
log.error("释放许可,当前可用许可数:{}",semaphore.availablePermits());
|
||||
}
|
||||
}, executor));
|
||||
try {
|
||||
NameplateTableKafkaDTO tmNameplateCorpus = JSON.parseObject(recordMessages, NameplateTableKafkaDTO.class);
|
||||
if (!"INSERT".equals(tmNameplateCorpus.getType()) || CollectionUtils.isEmpty(tmNameplateCorpus.getData())) {
|
||||
ack.acknowledge();
|
||||
return;
|
||||
}
|
||||
|
||||
log.info("nameplateKafkaConsumeracknowledge:{}",partitionId, offset);
|
||||
// 手动提交 offset
|
||||
} catch (Exception e) {
|
||||
log.info("nameplateKafkaConsumerFailed to process message: {}", e);
|
||||
}
|
||||
log.info("nameplateKafkaConsumer消息处理开始,耗时:{}", System.currentTimeMillis() - startTime);
|
||||
}
|
||||
// 使用CountDownLatch等待所有异步任务完成
|
||||
CountDownLatch latch = new CountDownLatch(tmNameplateCorpus.getData().size());
|
||||
AtomicBoolean hasError = new AtomicBoolean(false);
|
||||
|
||||
tmNameplateCorpus.getData().forEach(nameplate ->
|
||||
CompletableFuture.runAsync(() -> {
|
||||
try {
|
||||
if (!semaphore.tryAcquire(5, TimeUnit.SECONDS)) {
|
||||
log.warn("获取许可超时,跳过处理");
|
||||
hasError.set(true);
|
||||
return;
|
||||
}
|
||||
|
||||
tmNameplateCorpusService.processItem(nameplate);
|
||||
} catch (Exception e) {
|
||||
log.error("处理铭牌数据失败", e);
|
||||
hasError.set(true);
|
||||
} finally {
|
||||
semaphore.release();
|
||||
latch.countDown();
|
||||
}
|
||||
}, executor));
|
||||
|
||||
// 等待所有任务完成
|
||||
if (latch.await(30, TimeUnit.SECONDS)) {
|
||||
if (!hasError.get()) {
|
||||
ack.acknowledge(); // 只有所有任务成功才提交offset
|
||||
} else {
|
||||
log.error("部分任务失败,不提交offset");
|
||||
}
|
||||
} else {
|
||||
log.error("任务执行超时,不提交offset");
|
||||
}
|
||||
} catch (Exception e) {
|
||||
log.error("处理消息失败", e);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user