【发布时间】:2019-02-14 16:10:41
【问题描述】:
给定主题名称、分区号和偏移量,如何从主题中只读取一条记录?
在我的基于 Sprng Boot 的应用程序中,我使用 Kafka 导入业务数据。 导入记录被发送到 import_queue 并由一个或多个业务模块使用。即使消费者未能从记录中导入数据以继续从以下记录中导入数据,也会始终确认记录。
稍后用户(在他/她修复了一些相关的业务数据之后)可以决定重新发送一个或多个失败(但已确认)的导入记录。
每条记录的偏移量、分区号和主题名称都存储在我的应用程序内部的 SQL 数据库中。
从参考文档和一些 StackOverflow 问题中,我发现我必须:
- 设置容器(消费者/听众)
- 倒回(寻找)到所需的偏移量
- 读取一条记录
- 跳过读取剩余记录
这是从 kafka 主题中仅读取一条旧记录的唯一方法吗? 还是有更简单的解决方案?
解决方案
正如@Gary 所建议的:
ConsumerRecord<byte[], byte[]> read(String topic, int partition, long offset) {
Map<String, Object> configs = Map.of(
"bootstrap.servers", "localhost:9092",
"group.id", "incubator_retry",
"max.poll.records", 1);
DefaultKafkaConsumerFactory<byte[], byte[]> consumerFactory = new DefaultKafkaConsumerFactory<>(
configs, new ByteArrayDeserializer(), new ByteArrayDeserializer());
try (Consumer<byte[], byte[]> consumer = consumerFactory.createConsumer()) {
TopicPartition topicPartition = new TopicPartition(topic, partition);
consumer.assign(List.of(topicPartition));
consumer.seek(topicPartition, offset);
ConsumerRecords<byte[], byte[]> consumerRecords = consumer.poll(Duration.ofMillis(5000));
if (consumerRecords.isEmpty()) {
throw new RuntimeException(String.format("Timeout polling from topic %s partition %d at offset %d",
topicPartition.topic(), topicPartition.partition(), offset));
}
return consumerRecords.iterator().next();
}
}
【问题讨论】: