【发布时间】:2020-08-25 05:30:44
【问题描述】:
我正在使用 Quarkus 框架构建一个 Kafka 消费者,它将读取具有 3 个分区的主题。下面的代码 sn-p 正在工作,但基于日志,我只是用 3 个分区启动 1 个使用者。我现在的问题是如何在运行应用程序后生成 3 个消费者。
@Incoming("topic-1")
public CompletionStage<Void> onMessage(KafkaRecord<String, String> message) throws IOException {
LOG.info("Kafka order message with value = {} arrived from topic {} ", message.getPayload(),
message.getTopic());
//JsonObject event = new JsonObject(message.getPayload());
try {
if (true) {
LOG.info("Kafka message: " + message);
}
} catch (Exception e) {
e.printStackTrace();
}
return message.ack();
}
请查看示例日志:
INFO [org.apa.kaf.cli.con.int.ConsumerCoordinator] (vert.x-kafka-consumer-thread-0) [Consumer clientId=testconsumer, groupId=kafka-detection-consumer] 完成组分配在第 64 代:{testconsumer-bf6d314c-44e1-47b1-9439-fe4058951841=Assignment(partitions=[test_part-0, test_part-1, test_part-2])}
【问题讨论】:
标签: java apache-kafka reactive-programming kafka-consumer-api quarkus