【问题标题】:Messages produced before first consumer connected lost在第一个消费者连接之前产生的消息丢失
【发布时间】:2020-10-28 00:59:36
【问题描述】:

我在 kafka 中使用 kafka-topic.sh 创建了一个主题,并使用 java 客户端对其进行了测试:

kafka-topics.sh \
--create \
--zookeeper localhost:2181 \
--replication-factor 1 \
--partitions 2 \
--topic my-topic

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("my-topic"), new LoggingConsumerRebalanceListener(RandomStringUtils.randomAlphanumeric(3).toLowerCase()));
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(2000));
    for (ConsumerRecord<String, String> record : records)
        System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value());
    Thread.sleep(500);
}

Producer<String, String> producer = new KafkaProducer<>(props);
for (int i = 0; i < 10; i++) {
  String key = Integer.toString(i+1);
  String value = RandomStringUtils.randomAlphabetic(100);
  LOGGER.info("Sending message {}", key);
    producer.send(new ProducerRecord<String, String>("my-topic", key, value));
    Thread.sleep(100);
}
producer.close();    

生产者和消费者是我独立启动的独立代码块。

我有观察者,以下代码按顺序正常工作:

  • 设置主题
  • 运行消费者
  • 运行生产者
  • 运行生产者...

但是,按顺序:

  • 设置主题
  • 运行生产者 (1)
  • 运行消费者
  • 运行生产者

生产者第一次运行的消息丢失了。稍后,如果我停止消费者,运行生产者并运行消费者,我将收到所有消息。只有在第一个消费者订阅之前产生的消息会丢失。虽然我已经在命令行中明确创建了主题。

我在这里做错了什么?如何防止消息丢失?

【问题讨论】:

    标签: apache-kafka


    【解决方案1】:

    默认情况下,消费者将从最新的偏移量中读取。

    如果您运行“生产者 (1)”,然后启动消费者,它将忽略来自该生产者的消息并等待第二个生产者调用产生的新消息。

    可以通过配置 auto.offset.reset 更改从最新偏移量读取的行为。

    稍后,如果我停止消费者,运行生产者并运行消费者,我会收到所有消息

    发生这种情况是因为您的消费者有一个固定的消费者组(配置 group.id),并且默认设置 auto.offset.reset 不再有任何影响,因为该组已向 Kafka 注册,消费者将继续从主题中读取它停止的地方。

    最后,如果您不想在运行第二个序列时错过任何消息,请设置 auto.offset.reset=earliest 并定义一个新的唯一 group.id

    【讨论】:

      猜你喜欢
      • 2011-12-07
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-01-21
      • 2014-04-09
      • 2021-01-26
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多