【问题标题】:spring kafka offset increment even auto commit offset is set to falsespring kafka 偏移增量甚至自动提交偏移设置为 false
【发布时间】:2021-03-09 20:41:57
【问题描述】:

我正在尝试为在kafka 上收到的消息实施manual offset commit。我已将偏移提交设置为false,但偏移值不断增加。

不知道是什么原因。需要帮助解决问题。

下面是代码

application.yml

spring:
  application:
    name: kafka-consumer-sample
  resources:
    cache:
      period: 60m

kafka:
      bootstrapServers: localhost:9092
      options:
        enable:
          auto:
            commit: false

KafkaConfig.java

@Bean
    public ConsumerFactory<String, String> consumerFactory() {
        Map<String, Object> config = new HashMap<>();
        config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        config.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        config.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);

        return new DefaultKafkaConsumerFactory<>(config);
    }

 @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }

KafkaConsumer.java

@Service
public class KafkaConsumer {

    @KafkaListener(topics = "#{'${kafka-consumer.topics}'.split(',')}", groupId = "${kafka-consumer.groupId}")
    public void consume(ConsumerRecord<String, String> record) {

        System.out.println("Consumed Kafka Record: " + record);
        record.timestampType();
        System.out.println("record.timestamp() = " + record.timestamp());
        System.out.println("***********************************");
        System.out.println(record.timestamp());
        System.out.println("record.key() = " + record.key());
        System.out.println("Consumed String Message : " + record.value());
    }
}

输出如下

Consumed Kafka Record: ConsumerRecord(topic = test, partition = 0, offset = 31, CreateTime = 1573570989565, serialized key size = -1, serialized value size = 2, headers = RecordHeaders(headers = [], isReadOnly = false), key = null, value = 10)
record.timestamp() = 1573570989565
***********************************
1573570989565
record.key() = null
Consumed String Message : 10
Consumed Kafka Record: ConsumerRecord(topic = test, partition = 0, offset = 32, CreateTime = 1573570991535, serialized key size = -1, serialized value size = 2, headers = RecordHeaders(headers = [], isReadOnly = false), key = null, value = 11)
record.timestamp() = 1573570991535
***********************************
1573570991535
record.key() = null
Consumed String Message : 11

属性如下。

auto.commit.interval.ms = 100000000
auto.offset.reset = earliest
bootstrap.servers = [localhost:9092]
check.crcs = true
connections.max.idle.ms = 540000
enable.auto.commit = false
exclude.internal.topics = true
fetch.max.bytes = 52428800
fetch.max.wait.ms = 500
fetch.min.bytes = 1
group.id = mygroup
heartbeat.interval.ms = 3000

这是在我重新启动消费者之后。我希望早期的数据也会被打印出来。

我的理解正确吗? 请注意,我正在重新启动我的 springboot 应用程序,希望消息从第一个开始。并且我的 kafka 服务器和 zookeeper 没有终止。

【问题讨论】:

  • 我没有使用 poll() 方法。此外,即使每次后续轮询的偏移量都会自动增加,但它不会发生在服务器分区中。以及如何解决这个问题:(
  • 好的,你能展示一下 kafka 容器的配置吗?
  • 除了我已经提供的配置外,没有其他配置 :( 。您能指定您要查找的内容,以便我可以从我这边检查吗
  • @Deadpool 我已经用创建 ConcurrentKafkaListenerContainerFactory 的代码更新了我的问题。
  • 试试我的回答@Tushar Banne

标签: java apache-kafka kafka-consumer-api spring-kafka


【解决方案1】:

如果使用此属性ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG 禁用了auto 确认,那么您必须将容器级别的确认模式设置为MANUAL,并且不要提交offset,因为默认情况下它设置为BATCH.

@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory();
    factory.setConsumerFactory(consumerFactory());
    factory.getContainerProperties().setAckMode(AbstractMessageListenerContainer.AckMode.MANUAL);
    return factory;
}

因为在禁用自动确认时container 级别确认设置为BATCH

public void setAckMode(ContainerProperties.AckMode ackMode)

设置自动确认(在配置属性中)为假时使用的确认模式。

  1. RECORD:在将每条记录传递给侦听器后确认。
  2. BATCH:在从消费者收到的每批记录传递给侦听器后确认
  3. TIME:在此毫秒数后确认; (应该大于#setPollTimeout(long) pollTimeout。
  4. COUNT:至少收到此数量的记录后确认
  5. 手动:侦听器负责确认 - 使用 AcknowledgeingMessageListener。

参数:

ackMode - ContainerProperties.AckMode;默认 BATCH。

Committing Offsets

为提交偏移量提供了几个选项。如果 enable.auto.commit 消费者属性为 true,Kafka 会根据其配置自动提交偏移量。 如果为 false,则容器支持多个 AckMode 设置(在下一个列表中描述)。默认的 AckMode 是 BATCH。 从 2.3 版开始,除非在配置中明确设置,否则框架会将 enable.auto.commit 设置为 false。以前,如果未设置属性,则使用 Kafka 默认值 (true)。

如果您想始终从头开始阅读,则必须将此属性 auto.offset.reset 设置为 earliest

config.put(ConsumerConfig. AUTO_OFFSET_RESET_CONFIG, "earliest");

注意:确保groupId必须是新的,在kafka中没有任何偏移

【讨论】:

  • 我做了你所说的改变。开始消费。从 11 到 20 发送消息。停止服务并重新启动。发送了 21。我的期望是再次收到 11 到 20 的消息。但我没有得到 21 :(
  • 将该配置 AUTO_OFFSET_RESET_CONFIG 设置为 earliest @TusharBanne
  • 否 还是一样。偏移量不断增加。使用控制台上打印的最新属性更新帖子。
  • 尝试将消费者组 groupId = "${kafka-consumer.groupId} 更改为其他新组 @TusharBanne 我怀疑您使用的是较旧的消费者组,只需将其设置为 groupId = "test"
  • 是的。如果我更改组 ID,它会再次打印所有消息。那么这里的问题是什么?
猜你喜欢
  • 1970-01-01
  • 2021-11-10
  • 2018-03-27
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-01-13
相关资源
最近更新 更多