【问题标题】:In Spring-Kafka which method does the same thing like onPartitionsRevoked?在 Spring-Kafka 中,哪个方法可以像 onPartitionsRevoked 一样做同样的事情?
【发布时间】:2017-04-02 04:07:28
【问题描述】:

我知道在Sping-Kafka中我们有以下方法:

void registerSeekCallback(ConsumerSeekCallback 回调);

void onPartitionsAssigned(Map assignments, ConsumerSeekCallback 回调);

void onIdleContainer(Map assignments, ConsumerSeekCallback 回调);

但它与原生 ConsumerRebalanceListener 方法 onPartitionsRevoked 的作用相同?

"此方法将在重新平衡操作开始之前调用,并且 在消费者停止获取数据之后。建议抵消 应在此回调中提交给 Kafka 或自定义 偏移存储以防止重复数据。”

如果我想实现 ConsumerRebalanceListener,如何传递 KafkaConsumer 引用?我只看到 Spring-Kafka 的消费者。

=========更新======

嗨,Gary,当我将 RebalanceListener 添加到 ContainerProperties 中时。我可以看到这两种方法都被触发了。但是,我遇到了异常,说“提交无法完成,因为该组已经重新平衡并将分区分配给另一个成员。这意味着后续调用 poll() 之间的时间比配置的 max.poll 长。 interval.ms,这通常意味着轮询循环花费了太多时间处理消息“你知道吗?

===========更新2 ============

    public ConcurrentMessageListenerContainer<Integer, String> createContainer(
      ContainerProperties containerProps, IKafkaConsumer iKafkaConsumer) {

    Map<String, Object> props = consumerProps();

    DefaultKafkaConsumerFactory<Integer, String> cf = new DefaultKafkaConsumerFactory<Integer, String>(props);

    **RebalanceListner rebalanceListner = new RebalanceListner(cf.createConsumer());**

    CustomKafkaMessageListener ckml = new CustomKafkaMessageListener(iKafkaConsumer, rebalanceListner);

    CustomRecordFilter cff = new CustomRecordFilter();

    FilteringAcknowledgingMessageListenerAdapter faml = new FilteringAcknowledgingMessageListenerAdapter(ckml, cff, true);

    SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy();
    retryPolicy.setMaxAttempts(5);

    FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy();
    backOffPolicy.setBackOffPeriod(1500); // 1.5 seconds

    RetryTemplate rt = new RetryTemplate();
    rt.setBackOffPolicy(backOffPolicy);
    rt.setRetryPolicy(retryPolicy);
    rt.registerListener(ckml);
    RetryingAcknowledgingMessageListenerAdapter rml = new RetryingAcknowledgingMessageListenerAdapter(faml, rt);

    containerProps.setConsumerRebalanceListener(rebalanceListner);
    containerProps.setMessageListener(rml);
    containerProps.setAckMode(AbstractMessageListenerContainer.AckMode.MANUAL_IMMEDIATE);
    containerProps.setErrorHandler(ckml);
    containerProps.setAckOnError(false);
    ConcurrentMessageListenerContainer<Integer, String> container = new ConcurrentMessageListenerContainer<>(
        cf, containerProps);

    container.setConcurrency(1);
    return container;
  }

【问题讨论】:

    标签: apache-kafka spring-kafka


    【解决方案1】:

    您可以将RebalanceListener 添加到传递给构造函数的容器的ContainerProperties

    【讨论】:

    • 嗨,Gary,当我将 RebalanceListener 添加到 ContainerProperties 中时。我可以看到这两种方法都被触发了。但是,我遇到了异常,说“提交无法完成,因为该组已经重新平衡并将分区分配给另一个成员。这意味着后续调用 poll() 之间的时间比配置的 max.poll 长。 interval.ms,这通常意味着轮询循环花费了太多时间处理消息“你有什么想法吗?
    • 我假设您使用的是手动确认模式之一。您不必自己确认提交 - 容器将在调用您的侦听器之前代表您这样做 - 在重新平衡之前未完成的任何待处理确认都将被提交。如果您没有看到该行为,请分享您的配置。请参阅the code here - 在我们提交未完成的偏移量后,您的侦听器在第 409 行被调用。
    • 是的。加里,我设置为 auto.commit = false。我的监听器正在实现 AcknowledgeingMessageListener 并且 AckMode 设置为 AbstractMessageListenerContainer.AckMode.MANUAL_IMMEDIATE。
    • Gary,我只使用原生的 Kafka Consumer API,并且能够实现 RebalanceListener。但使用 Spring-Kafka 我仍然得到异常。我将在我的帖子中更新我的配置。请看一下。
    • 我想你可能发现了一个错误。容器(以及一般的 spring-kafka)试图让开发人员的事情变得更容易。为此,当发生重新平衡时,容器会提交任何待处理的确认(您已经调用了acknowledge(),但它们尚未提交);然后它会调用您的侦听器,但不会检查您的侦听器是否添加了更多确认 - 您可以显示您的重新平衡侦听器代码吗?如果是这样,请在github 上打开一个问题。
    猜你喜欢
    • 2020-08-05
    • 1970-01-01
    • 2010-09-18
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-06-26
    相关资源
    最近更新 更多