修改kafka配置

This commit is contained in:
zren25
2025-05-16 13:28:40 +08:00
parent 0a980c78f5
commit 3865912527
2 changed files with 11 additions and 8 deletions

View File

@@ -65,10 +65,10 @@ public class KafkaConfig {
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, valueDeserializer); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, valueDeserializer);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest"); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 120000); // 会话1分钟 props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 45000); // 会话1分钟
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300000); // 可选:允许更长的消费间隔 props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 600000); // 可选:允许更长的消费间隔
props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 10000); props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 15000);
/** props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 20); // 单次最多拉取的消息数 **/ props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 20); // 单次最多拉取的消息数
ConcurrentKafkaListenerContainerFactory<String, String> factory = ConcurrentKafkaListenerContainerFactory<String, String> factory =
new ConcurrentKafkaListenerContainerFactory<>(); new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(props)); factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(props));

View File

@@ -70,12 +70,15 @@ 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 {
// 手动提交 offset
if (ack != null) {
log.info("nameplateKafkaConsumerack.acknowledge"); log.info("nameplateKafkaConsumerack.acknowledge");
ack.acknowledge(); try {
ack.acknowledge(); // 确保无论成功与否都提交 offset
} catch (IllegalStateException e) {
log.warn("Offset 已提交,跳过重复提交");
} }
}
}else { }else {
ack.acknowledge(); // 空消息直接跳过 ack.acknowledge(); // 空消息直接跳过
} }