【发布时间】: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