修改kafka配置
This commit is contained in:
@@ -65,8 +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, 60000); // 会话1分钟
|
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 120000); // 会话1分钟
|
||||||
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300000); // 可选:允许更长的消费间隔
|
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300000); // 可选:允许更长的消费间隔
|
||||||
|
props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 10000);
|
||||||
|
/** 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));
|
||||||
|
|||||||
@@ -53,26 +53,6 @@ public class CorpusProcessKafkaProducer {
|
|||||||
@Resource
|
@Resource
|
||||||
private RocketMQTemplate rocketMqTemplate;
|
private RocketMQTemplate rocketMqTemplate;
|
||||||
|
|
||||||
// @KafkaListener(
|
|
||||||
// topicPartitions = @TopicPartition(
|
|
||||||
// topic = "${spring.kafka.topic}",
|
|
||||||
// partitionOffsets = {
|
|
||||||
// @PartitionOffset(partition = "0", initialOffset = "2792520"),
|
|
||||||
// @PartitionOffset(partition = "1", initialOffset = "2596153"),
|
|
||||||
// @PartitionOffset(partition = "2", initialOffset = "2536889"),
|
|
||||||
// @PartitionOffset(partition = "3", initialOffset = "2782173"),
|
|
||||||
// @PartitionOffset(partition = "4", initialOffset = "2616677"),
|
|
||||||
// @PartitionOffset(partition = "5", initialOffset = "2529585"),
|
|
||||||
// @PartitionOffset(partition = "6", initialOffset = "2761950"),
|
|
||||||
// @PartitionOffset(partition = "7", initialOffset = "2581957"),
|
|
||||||
// @PartitionOffset(partition = "8", initialOffset = "2530225"),
|
|
||||||
// @PartitionOffset(partition = "9", initialOffset = "2803179"),
|
|
||||||
// @PartitionOffset(partition = "10", initialOffset = "2599129"),
|
|
||||||
// @PartitionOffset(partition = "11", initialOffset = "2546277")
|
|
||||||
// }
|
|
||||||
// ),
|
|
||||||
// groupId = "${spring.kafka.group}"
|
|
||||||
// )
|
|
||||||
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}")
|
@KafkaListener(topics = "${spring.kafka.topic}", groupId = "${spring.kafka.group}")
|
||||||
public void listen(List<ConsumerRecord<String, Object>> recordMessages) {
|
public void listen(List<ConsumerRecord<String, Object>> recordMessages) {
|
||||||
long startTime = System.currentTimeMillis();
|
long startTime = System.currentTimeMillis();
|
||||||
|
|||||||
@@ -13,6 +13,8 @@ import org.springframework.beans.factory.annotation.Value;
|
|||||||
import org.springframework.cloud.context.config.annotation.RefreshScope;
|
import org.springframework.cloud.context.config.annotation.RefreshScope;
|
||||||
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.messaging.handler.annotation.Header;
|
||||||
import org.springframework.stereotype.Component;
|
import org.springframework.stereotype.Component;
|
||||||
import org.springframework.web.bind.annotation.RestController;
|
import org.springframework.web.bind.annotation.RestController;
|
||||||
|
|
||||||
@@ -48,7 +50,9 @@ public class NameplateKafkaConsumer {
|
|||||||
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.OFFSET) Long offset) {
|
||||||
long startTime = System.currentTimeMillis();
|
long startTime = System.currentTimeMillis();
|
||||||
|
|
||||||
log.info("nameplateKafkaConsumer 当前线程: {}, 线程ID: {}", Thread.currentThread().getName(), Thread.currentThread().getId());
|
log.info("nameplateKafkaConsumer 当前线程: {}, 线程ID: {}", Thread.currentThread().getName(), Thread.currentThread().getId());
|
||||||
@@ -64,7 +68,8 @@ public class NameplateKafkaConsumer {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
} catch (Exception e) {
|
} catch (Exception e) {
|
||||||
log.error("nameplateKafkaConsumer铭牌解析处理出错: {}", e.getMessage());
|
log.error("nameplateKafkaConsumer铭牌解析处理出错 {}:{},ex:{}",
|
||||||
|
partitionId, offset, e.getMessage());
|
||||||
}
|
}
|
||||||
// 手动提交 offset
|
// 手动提交 offset
|
||||||
if (ack != null) {
|
if (ack != null) {
|
||||||
|
|||||||
Reference in New Issue
Block a user