【发布时间】:2020-03-13 09:30:01
【问题描述】:
我总是担心卡夫卡 1. 重复 2. 缺失记录
我在春季 Kafka 2.2.2.RELEASE 中进行了以下更改以解决上述问题。 有人可以确认这是否正确。
public ConsumerFactory<String, String> consumerFactory() {
Map<String, Object> props = new HashMap<>();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "test-consumer-group");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 20);
props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 600000);
props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 1000);
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 60000);
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
return new DefaultKafkaConsumerFactory<>(props);
}
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
factory.getContainerProperties().setCommitLogLevel(LogIfLevelEnabled.Level.INFO);
factory.getContainerProperties().setAckOnError(false);
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
return factory;
另外请确认,我是否需要在 ConcurrentKafkaListenerContainerFactory 中实现 ConsumerRebalanceListener。
如果需要如何实现。
factory.getContainerProperties().setConsumerRebalanceListener(new ConsumerRebalanceListener() {
@Override
public void onPartitionsRevoked(Collection<TopicPartition> collection) {
//TODO code
}
@Override
public void onPartitionsAssigned(Collection<TopicPartition> collection) {
//TODO code
}
});
【问题讨论】:
-
如果您能提供一些例子说明我们将此记录重复或在总记录中缺少某些内容会有所帮助
-
添加下面的配置后,我没有设置重复记录。 props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 20); props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 600000); props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 1000); props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 60000);在此之前,我的日志中有重复记录。那么这样可以吗?但目前我收到一些记录丢失。所以我正在分析这个
-
假设有数百万条记录进入kafka消费者,以上配置是否可以。 props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 20); props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 600000);还是我需要改变什么?
标签: spring-kafka