【发布时间】:2020-04-15 00:09:12
【问题描述】:
我正在使用 spring-kafka '2.1.7.RELEASE' 并且我试图了解 max.poll.interval.ms 如何将 AckMode 作为 BATCH 并将 enable.auto.commit 作为 'false' 工作。这是我的消费者设置。
public Map<String, Object> setConsumerConfigs() {
Map<String, Object> configs = = new HashMap<>();
configs.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
configs.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
configs.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, "400000");
configs.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer2.class);
configs.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer2.class);
configs.put(ErrorHandlingDeserializer2.KEY_DESERIALIZER_CLASS, stringDeserializerClass);
configs.put(ErrorHandlingDeserializer2.VALUE_DESERIALIZER_CLASS, kafkaAvroDeserializerClass.getName());
configs.setPartitionAssignmentStrategyConfig(Collections.singletonList(RoundRobinAssignor.class));
// Set this to true so that you will have consumer record value coming as your pre-defined contract instead of a generic record
sapphireKafkaConsumerConfig.setSpecificAvroReader("true");
}
这是我的出厂设置
@Bean
public <K,V> ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(getConsumerConfigs));
factory.getContainerProperties().setMissingTopicsFatal(false);
factory.getContainerProperties().setAckMode(AckMode.BATCH);
factory.setErrorHandler(myCustomKafkaSeekToCurrentErrorHandler);
factory.setRetryTemplate(retryTemplate());
factory.setRecoveryCallback(myCustomKafkaRecoveryCallback);
factory.setStatefulRetry(true);
return factory;
}
public RetryTemplate retryTemplate() {
RetryTemplate retryTemplate = new RetryTemplate();
retryTemplate.setListeners(new RetryListener[]{myCustomKafkaRetryListener});
retryTemplate.setRetryPolicy(myCustomKafkaConsumerRetryPolicy);
FixedBackOffPolicy backOff = new FixedBackOffPolicy();
backOff.setBackOffPeriod(1000);
retryTemplate.setBackOffPolicy(backOff);
return retryTemplate;
}
这是我的消费者,我添加了 2 分钟的延迟
@KafkaListener(topics = TestConsumerConstants.CONSUMER_LONGRUNNING_RECORDS_PROCESSSING_TEST_TOPIC
, clientIdPrefix = "CONSUMER_LONGRUNNING_RECORDS_PROCESSSING"
, groupId = "kafka-lib-comp-test-consumers")
public void consumeLongRunningRecord(ConsumerRecord message) throws InterruptedException {
System.out.println(String.format("\n \n Received message at %s offset %s of partition %s of topic %s with key %s \n\n", DateTime.now(),
message.offset(), message.partition(), message.topic(), message.key()));
TimeUnit.MINUTES.sleep(2);
System.out.println(String.format("\n \n Processing done for the message at %s offset %s of partition %s of topic %s with key %s \n\n", DateTime.now(),
message.offset(), message.partition(), message.topic(), message.key()));
}
现在,我发布了 5 条消息,并观察到它处理了所有记录,没有任何问题。但是,如果我将 AckMode 设置为 RECORD,它会在处理第 4 条消息后提交偏移量时抛出错误,然后处理相同的消息两次(这是预期的)。
根据 spring-kafka 文档,当 poll() 返回的所有记录都被处理后,AckMode = BATCH 将提交偏移量。
现在,我的问题是,AckMode 如何在通过 max.poll.interval.ms 后改变行为而不引起重新平衡?请帮我理解。
Caused by: org.apache.kafka.clients.consumer.CommitFailedException: Commit cannot be completed since the group has already rebalanced and assigned the partitions to another member. This means that the time between subsequent calls to poll() was longer than the configured max.poll.interval.ms, which typically implies that the poll loop is spending too much time message processing. You can address this either by increasing the session timeout or by reducing the maximum size of batches returned in poll() with max.poll.records.
【问题讨论】:
标签: spring-kafka