【问题标题】:How to set multiple consumers in the same topic in Kafka using Quarkus如何使用 Quarkus 在 Kafka 的同一主题中设置多个消费者
【发布时间】: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


    【解决方案1】:

    如果您在容器平台(Docker、K8S ...)上运行您的应用程序,那么您可以水平扩展您的服务;否则,请使用不同的端口再次运行您的应用程序。

    Kafka客户端启动时,会被分配到某个partition,同一个客户端不能消费多个topic-partition。

    【讨论】:

    • 这可以是一个选项,但我想要在服务中增加线程或其他选项,我可以在服务中编写代码。
    • 我没试过,但也许你可以配置另一个指向同一个Kafka主题的通道名称,并创建两个具有相同代码但通道名称不同的方法,我相信应该这样做。
    • 谢谢!这是我的最后一个选择,但它会有一个脏代码 - 我正在寻找动态配置,并将根据该配置生成线程/消费者。
    • 还要确保所有消费者使用相同的消费者组 ID。
    猜你喜欢
    • 2022-01-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-01-26
    • 1970-01-01
    • 2019-09-26
    • 1970-01-01
    • 2017-12-23
    相关资源
    最近更新 更多