修改kafka次数限制
This commit is contained in:
@@ -14,6 +14,8 @@ import org.springframework.kafka.listener.ContainerProperties;
|
|||||||
|
|
||||||
import java.util.HashMap;
|
import java.util.HashMap;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
import java.util.concurrent.ExecutorService;
|
||||||
|
import java.util.concurrent.Executors;
|
||||||
|
|
||||||
@Configuration
|
@Configuration
|
||||||
@RefreshScope
|
@RefreshScope
|
||||||
@@ -68,7 +70,7 @@ public class KafkaConfig {
|
|||||||
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 45000); // 会话1分钟
|
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 45000); // 会话1分钟
|
||||||
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 600000); // 可选:允许更长的消费间隔
|
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 600000); // 可选:允许更长的消费间隔
|
||||||
props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 15000);
|
props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 15000);
|
||||||
/** props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 20); // 单次最多拉取的消息数 **/
|
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 1); // 单次最多拉取的消息数 **/
|
||||||
ConcurrentKafkaListenerContainerFactory<String, String> factory =
|
ConcurrentKafkaListenerContainerFactory<String, String> factory =
|
||||||
new ConcurrentKafkaListenerContainerFactory<>();
|
new ConcurrentKafkaListenerContainerFactory<>();
|
||||||
factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(props));
|
factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(props));
|
||||||
@@ -77,4 +79,10 @@ public class KafkaConfig {
|
|||||||
return factory;
|
return factory;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Bean("kafkaTaskExecutor")
|
||||||
|
public ExecutorService kafkaTaskExecutor() {
|
||||||
|
int optimalThreadPoolSize = Runtime.getRuntime().availableProcessors() + 1;
|
||||||
|
return Executors.newFixedThreadPool(optimalThreadPoolSize);
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
@@ -100,7 +100,7 @@ public class CorpusFailJob {
|
|||||||
queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisErrors.getAiAnalysisRequestId());
|
queryWrapper.eq(AiAnalysisRequestLogs::getAiAnalysisRequestId, aiAnalysisErrors.getAiAnalysisRequestId());
|
||||||
AiAnalysisRequestLogs oldAiAnalysisRequestLogs = aiAnalysisRequestLogsMapper.selectOne(queryWrapper);
|
AiAnalysisRequestLogs oldAiAnalysisRequestLogs = aiAnalysisRequestLogsMapper.selectOne(queryWrapper);
|
||||||
|
|
||||||
if (null != oldAiAnalysisRequestLogs && StringUtils.isNotBlank(oldAiAnalysisRequestLogs.getBusinessResponse())) {
|
if (null != oldAiAnalysisRequestLogs && StringUtils.isBlank(oldAiAnalysisRequestLogs.getBusinessResponse())) {
|
||||||
DiFyReq diFyReq = JSONObject.parseObject(oldAiAnalysisRequestLogs.getDifyRequest(), DiFyReq.class);
|
DiFyReq diFyReq = JSONObject.parseObject(oldAiAnalysisRequestLogs.getDifyRequest(), DiFyReq.class);
|
||||||
CorpusReportDTO corpusReportDTO = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class);
|
CorpusReportDTO corpusReportDTO = JSONObject.parseObject(oldAiAnalysisRequestLogs.getBusinessRequest(), CorpusReportDTO.class);
|
||||||
Map<String, String> ltoMap = new HashMap<>();
|
Map<String, String> ltoMap = new HashMap<>();
|
||||||
@@ -141,7 +141,7 @@ public class CorpusFailJob {
|
|||||||
|
|
||||||
}
|
}
|
||||||
} else if(null != corpusReportDTO && corpusReportDTO.getAnalysisScene() == 3){
|
} else if(null != corpusReportDTO && corpusReportDTO.getAnalysisScene() == 3){
|
||||||
log.info(" corpusFailTask 铭牌数据重试:{}", corpusReportDTO.getAnalysisScene());
|
log.info(" corpusFailTask 铭牌数据重试:{},aiId:{}", corpusReportDTO.getCustomerFlowId(), oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
|
||||||
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyReq, oldAiAnalysisRequestLogs.getAiAnalysisRequestType(), JSONObject.toJSONString(corpusReportDTO),oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
|
JSONObject execDifyFlow = diFyService.executeDifyFlow(diFyReq, oldAiAnalysisRequestLogs.getAiAnalysisRequestType(), JSONObject.toJSONString(corpusReportDTO),oldAiAnalysisRequestLogs.getAiAnalysisRequestId());
|
||||||
List<TmNameplateCorpus> nameplateCorpusList = tmNameplateCorpusService.queryTelephoneCorpusByCustomerFlowId( Arrays.asList(corpusReportDTO.getCustomerFlowId()));
|
List<TmNameplateCorpus> nameplateCorpusList = tmNameplateCorpusService.queryTelephoneCorpusByCustomerFlowId( Arrays.asList(corpusReportDTO.getCustomerFlowId()));
|
||||||
if(CollectionUtils.isNotEmpty(nameplateCorpusList)) {
|
if(CollectionUtils.isNotEmpty(nameplateCorpusList)) {
|
||||||
|
|||||||
@@ -0,0 +1,54 @@
|
|||||||
|
package com.volvo.ai.analytic.center.mq;
|
||||||
|
|
||||||
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
|
import org.springframework.stereotype.Component;
|
||||||
|
|
||||||
|
import java.util.concurrent.Executors;
|
||||||
|
import java.util.concurrent.ScheduledExecutorService;
|
||||||
|
import java.util.concurrent.TimeUnit;
|
||||||
|
import java.util.concurrent.atomic.AtomicInteger;
|
||||||
|
|
||||||
|
@Component
|
||||||
|
public class ConsumerStateManager {
|
||||||
|
// 使用 ThreadLocal 存储每个线程的 Consumer 实例
|
||||||
|
private final ThreadLocal<Consumer<?, ?>> consumerThreadLocal = new ThreadLocal<>();
|
||||||
|
private final AtomicInteger messageCount = new AtomicInteger(0);
|
||||||
|
private volatile boolean paused = false;
|
||||||
|
private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
|
||||||
|
|
||||||
|
// 绑定当前线程的 Consumer
|
||||||
|
public void bindConsumer(Consumer<?, ?> consumer) {
|
||||||
|
consumerThreadLocal.set(consumer);
|
||||||
|
}
|
||||||
|
|
||||||
|
// 获取当前线程绑定的 Consumer
|
||||||
|
public Consumer<?, ?> getKafkaConsumer() {
|
||||||
|
return consumerThreadLocal.get();
|
||||||
|
}
|
||||||
|
|
||||||
|
// 更新计数器并检查是否需要暂停
|
||||||
|
public void checkAndPause() {
|
||||||
|
if (messageCount.incrementAndGet() >= 30 && !paused) {
|
||||||
|
paused = true;
|
||||||
|
Consumer<?, ?> consumer = getKafkaConsumer();
|
||||||
|
if (consumer != null) {
|
||||||
|
consumer.pause(consumer.assignment());
|
||||||
|
scheduler.schedule(this::resume, 2, TimeUnit.MINUTES);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 恢复消费
|
||||||
|
private void resume() {
|
||||||
|
paused = false;
|
||||||
|
messageCount.set(0);
|
||||||
|
Consumer<?, ?> consumer = getKafkaConsumer();
|
||||||
|
if (consumer != null) {
|
||||||
|
consumer.resume(consumer.assignment());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
public boolean isPaused() {
|
||||||
|
return paused;
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -6,11 +6,14 @@ import com.fasterxml.jackson.databind.ObjectMapper;
|
|||||||
import com.volvo.ai.analytic.center.dto.corpus.NameplateTableKafkaDTO;
|
import com.volvo.ai.analytic.center.dto.corpus.NameplateTableKafkaDTO;
|
||||||
import com.volvo.ai.analytic.center.service.TmNameplateCorpusService;
|
import com.volvo.ai.analytic.center.service.TmNameplateCorpusService;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
|
import org.apache.commons.collections.CollectionUtils;
|
||||||
import org.apache.commons.lang3.StringUtils;
|
import org.apache.commons.lang3.StringUtils;
|
||||||
|
import org.apache.kafka.clients.consumer.Consumer;
|
||||||
import org.apache.rocketmq.spring.core.RocketMQTemplate;
|
import org.apache.rocketmq.spring.core.RocketMQTemplate;
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.beans.factory.annotation.Value;
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
import org.springframework.cloud.context.config.annotation.RefreshScope;
|
import org.springframework.cloud.context.config.annotation.RefreshScope;
|
||||||
|
import org.springframework.context.annotation.Bean;
|
||||||
import org.springframework.kafka.annotation.KafkaListener;
|
import org.springframework.kafka.annotation.KafkaListener;
|
||||||
import org.springframework.kafka.support.Acknowledgment;
|
import org.springframework.kafka.support.Acknowledgment;
|
||||||
import org.springframework.kafka.support.KafkaHeaders;
|
import org.springframework.kafka.support.KafkaHeaders;
|
||||||
@@ -19,9 +22,8 @@ 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.CompletableFuture;
|
import java.util.concurrent.*;
|
||||||
import java.util.concurrent.ExecutorService;
|
import java.util.concurrent.atomic.AtomicInteger;
|
||||||
import java.util.concurrent.Executors;
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* @ClassName NameplateKafkaConsumer
|
* @ClassName NameplateKafkaConsumer
|
||||||
@@ -42,56 +44,84 @@ public class NameplateKafkaConsumer {
|
|||||||
|
|
||||||
private final ObjectMapper objectMapper = new ObjectMapper();
|
private final ObjectMapper objectMapper = new ObjectMapper();
|
||||||
|
|
||||||
@Value("${rocketmq.producer.corpus.dcctopic}")
|
// @Value("${analyticCenterKafka.consumer.threshold}")
|
||||||
private String dccMqTipic;
|
// private String threshold;
|
||||||
|
|
||||||
@Resource
|
@Resource
|
||||||
private RocketMQTemplate rocketMqTemplate;
|
private RocketMQTemplate rocketMqTemplate;
|
||||||
|
|
||||||
private int messageCount = 0;
|
private final AtomicInteger messageCount = new AtomicInteger(0);
|
||||||
@KafkaListener(topics = "${analyticCenterKafka.consumer.topic}",
|
// 批次阈值(例如每30条提交一次offset)
|
||||||
|
private static final int THRESHOLD = 20;
|
||||||
|
|
||||||
|
// 最大等待时间(2分钟)
|
||||||
|
private static final long MAX_WAIT_TIME = TimeUnit.MINUTES.toMillis(2);
|
||||||
|
@Autowired
|
||||||
|
@Qualifier("kafkaTaskExecutor")
|
||||||
|
public ExecutorService kafkaTaskExecutor;
|
||||||
|
|
||||||
|
@Autowired
|
||||||
|
public ConsumerStateManager consumerStateManager;
|
||||||
|
@KafkaListener(topics = "${analyticCenterKafka.consumer.topic}",
|
||||||
groupId = "${analyticCenterKafka.consumer.group}" ,
|
groupId = "${analyticCenterKafka.consumer.group}" ,
|
||||||
containerFactory = "analyticCenterConsumerFactory",
|
containerFactory = "analyticCenterConsumerFactory",
|
||||||
concurrency = "3")
|
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) {
|
@Header(KafkaHeaders.OFFSET) Long offset,
|
||||||
|
Consumer<?, ?> consumer ) throws InterruptedException {
|
||||||
long startTime = System.currentTimeMillis();
|
long startTime = System.currentTimeMillis();
|
||||||
|
log.info("nameplateKafkaConsumer 当前线程: {}, 线程ID: {},计数:{}", Thread.currentThread().getName(), Thread.currentThread().getId());
|
||||||
messageCount++;
|
|
||||||
log.info("nameplateKafkaConsumer 当前线程: {}, 线程ID: {},计数:{}", Thread.currentThread().getName(), Thread.currentThread().getId(),messageCount);
|
|
||||||
log.info("nameplateKafkaConsumerMessage: {}", recordMessages);
|
log.info("nameplateKafkaConsumerMessage: {}", recordMessages);
|
||||||
if(StringUtils.isNotEmpty(recordMessages)){
|
// 初始化绑定 Consumer
|
||||||
|
|
||||||
int optimalThreadPoolSize = Runtime.getRuntime().availableProcessors() + 1;
|
consumerStateManager.bindConsumer(consumer);
|
||||||
log.info("获取的线程数:{}",optimalThreadPoolSize);
|
|
||||||
// 创建线程池
|
// 检查是否处于暂停状态
|
||||||
ExecutorService executor = Executors.newFixedThreadPool(optimalThreadPoolSize); // 根据需求调整线程池大小
|
if (consumerStateManager.isPaused()) {
|
||||||
NameplateTableKafkaDTO tmNameplateCorpus = JSON.parseObject(recordMessages, NameplateTableKafkaDTO.class);
|
log.info("当前处于暂停状态,忽略消息: partition={}, offset={}", partitionId, offset);
|
||||||
log.info("nameplateKafkaConsumerParseType: {},size:{}", tmNameplateCorpus.getType(), tmNameplateCorpus.getData().size());
|
return; // 不提交偏移量,等待恢复后重新拉取
|
||||||
if(tmNameplateCorpus.getType().equals("INSERT")){
|
|
||||||
tmNameplateCorpus.getData().forEach(nameplate ->
|
|
||||||
CompletableFuture.runAsync(() -> {
|
|
||||||
try {
|
|
||||||
log.info("nameplateKafkaConsumerCustomerFlowId: {}",nameplate.getCustomerFlowId());
|
|
||||||
tmNameplateCorpusService.processItem(nameplate);
|
|
||||||
} catch (Exception e) {
|
|
||||||
log.error("nameplateKafkaConsumer铭牌解析处理出错 {}:{},ex:{}",
|
|
||||||
partitionId, offset, e.getMessage());
|
|
||||||
}
|
|
||||||
}, executor));
|
|
||||||
try {
|
|
||||||
log.info("nameplateKafkaConsumerack.acknowledge:{}");
|
|
||||||
ack.acknowledge(); // 确保无论成功与否都提交 offset
|
|
||||||
} catch (IllegalStateException e) {
|
|
||||||
log.warn("Offset 已提交,跳过重复提交");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}else {
|
|
||||||
ack.acknowledge(); // 空消息直接跳过
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (StringUtils.isEmpty(recordMessages)) {
|
||||||
|
ack.acknowledge();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
NameplateTableKafkaDTO tmNameplateCorpus = JSON.parseObject(recordMessages, NameplateTableKafkaDTO.class);
|
||||||
|
if (!"INSERT".equals(tmNameplateCorpus.getType()) || CollectionUtils.isEmpty(tmNameplateCorpus.getData())) {
|
||||||
|
ack.acknowledge();
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
|
||||||
|
messageCount.addAndGet(1);
|
||||||
|
if (messageCount.get() >= THRESHOLD) {
|
||||||
|
log.info("nameplateKafkaConsumer 消费数量:{}", messageCount.get());
|
||||||
|
Thread.sleep(MAX_WAIT_TIME); // 让线程休眠指定的时间
|
||||||
|
messageCount.set(0); // 重置计数器
|
||||||
|
}
|
||||||
|
|
||||||
|
log.info("nameplateKafkaConsumerParseType: {},size:{}", tmNameplateCorpus.getType(), tmNameplateCorpus.getData().size());
|
||||||
|
tmNameplateCorpus.getData().forEach(nameplate ->
|
||||||
|
CompletableFuture.runAsync(() -> {
|
||||||
|
try {
|
||||||
|
log.info("nameplateKafkaConsumerCustomerFlowId: {}",nameplate.getCustomerFlowId());
|
||||||
|
tmNameplateCorpusService.processItem(nameplate);
|
||||||
|
} catch (Exception e) {
|
||||||
|
log.error("nameplateKafkaConsumer铭牌解析处理出错 {}:{},ex:{}",
|
||||||
|
partitionId, offset, e.getMessage());
|
||||||
|
}
|
||||||
|
}, kafkaTaskExecutor));
|
||||||
|
|
||||||
|
try {
|
||||||
|
ack.acknowledge(); // 确保无论成功与否都提交 offset
|
||||||
|
consumerStateManager.checkAndPause(); // 更新计数器并检查暂停
|
||||||
|
} catch (Exception e) {
|
||||||
|
log.error("任务执行异常: {}", e.getMessage());
|
||||||
|
Thread.currentThread().interrupt();
|
||||||
|
}
|
||||||
log.info("nameplateKafkaConsumer消息处理完成,耗时:{}", System.currentTimeMillis() - startTime);
|
log.info("nameplateKafkaConsumer消息处理完成,耗时:{}", System.currentTimeMillis() - startTime);
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user