【问题标题】:How to call kafkaconsumer api from partition assignor' s implementation如何从分区分配器的实现中调用 kafkaconsumer api
【发布时间】:2021-01-18 05:40:38
【问题描述】:

我通过在我的 Spring Boot 应用程序中实现 RangeAssignor 实现了我自己的分区分配策略。 我已经覆盖了它的 subscriptionUserData 方法并添加了一些用户数据。每当这些数据发生变化时,我想通过调用以下 kafkaConsumer 的 api 来触发分区重新平衡 kafkaconsumer apis enforce rebalance

我不确定如何获取 kafka 消费者的对象并调用此 api。 请推荐

【问题讨论】:

    标签: spring-boot apache-kafka


    【解决方案1】:

    你可以调用 consumer.wakeup() 函数

    consumer.wakeup() 是唯一可以安全地从不同线程调用的消费者方法。调用 wakeup 会导致 poll() 以 WakeupException 退出,或者如果在线程未等待 poll 时调用了 consumer.wakeup(),则在调用 poll() 时将在下一次迭代中抛出异常。 WakeupException 不需要处理,但在退出线程之前,必须调用 consumer.close()。如果需要,关闭消费者将提交偏移量,并将向组协调器发送一条消息,告知消费者将离开组。消费者协调器将立即触发重新平衡

      Runtime.getRuntime().addShutdownHook(new Thread() {
            public void run() {
                System.out.println("Starting exit...");
                consumer.wakeup();   **//1**
                try {
                    mainThread.join();
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
    } });
        ...
        Duration timeout = Duration.ofMillis(100);
        try {
            // looping until ctrl-c, the shutdown hook will cleanup on exit
            while (true) {
                ConsumerRecords<String, String> records =
                    movingAvg.consumer.poll(timeout);
                System.out.println(System.currentTimeMillis() +
                    "--  waiting for data...");
                for (ConsumerRecord<String, String> record : records) {
                    System.out.printf("offset = %d, key = %s, value = %s\n",
                        record.offset(), record.key(), record.value());
                }
                for (TopicPartition tp: consumer.assignment())
                    System.out.println("Committing offset at position:" +
                        consumer.position(tp));
                    movingAvg.consumer.commitSync();
            }
        } catch (WakeupException e) {
            // ignore for shutdown. **//2**
        } finally {
            consumer.close(); **//3**
            System.out.println("Closed consumer and we are done");
        }
    
    1. ShutdownHook 在单独的线程中运行,因此我们可以采取的唯一安全操作是调用 wakeup 以跳出轮询循环。
    2. 另一个调用wakeup 的线程将导致poll 抛出WakeupException。您需要捕获异常以确保您的应用程序不会意外退出,但无需对其执行任何操作。
    3. 在退出消费者之前,请确保将其关闭干净。

    完整示例在:

    https://github.com/gwenshap/kafka-examples/blob/master/SimpleMovingAvg/src/main/java/com/shapira/examples/newconsumer/simplemovingavg/SimpleMovingAvgNewConsumer.java

    【讨论】:

    • 感谢您的回答,但这不是我需要的,我需要一种在我的 custompartition 分配器类中获取 kafkaconsumer 对象的方法
    • 你可以调用 consumer.wakeup();来自您的自定义分区分配器类。这将通过您的 kafka 消费者对象开始重新平衡
    • 嗨 Sudhir,我也有类似的问题,您现在对如何解决问题有什么建议吗(当自定义用户数据被修改时触发分区重新平衡)?顺便说一句,您是否考虑过 github.com/grantneale/kafka-lag-based-assignor/blob/master/src/… 中所做的元数据使用者?或者您找到了更好的解决方案?谢谢?
    猜你喜欢
    • 1970-01-01
    • 2017-04-21
    • 2017-04-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-07-09
    相关资源
    最近更新 更多