【问题标题】:Kafka consumer receives message if set group_id to None, but it doesn't receive any message if not None?如果将group_id设置为None,Kafka消费者会收到消息,但如果不是None,它不会收到任何消息?
【发布时间】:2020-11-07 00:32:15
【问题描述】:

我有以下 Kafka 消费者,如果将 group_id 分配给 None 效果很好 - 它收到了所有历史消息和我新测试的消息。

consumer = KafkaConsumer(
        topic,
        bootstrap_servers=bootstrap_servers,
        auto_offset_reset=auto_offset_reset,
        enable_auto_commit=enable_auto_commit,
        group_id=group_id,
        value_deserializer=lambda x: json.loads(x.decode('utf-8'))
    )

for m in consumer:

但是,如果我将group_id 设置为某个值,它不会收到任何信息。我尝试运行测试生产者发送新消息,但没有收到任何消息。

消费者控制台确实显示以下消息:

2020-11-07 00:56:01 INFO ThreadPoolExecutor-0_0 base.py(重新)加入组 my_group 2020-11-07 00:56:07 INFO ThreadPoolExecutor-0_0 base.py 成功加入组 my_group 与第 497 代 2020-11-07 00:56:07 INFO ThreadPoolExecutor-0_0 subscription_state.py 更新了分区分配:[] 2020-11-07 00:56:07 INFO ThreadPoolExecutor-0_0 consumer.py 为组 my_group 设置新分配的分区 set()

【问题讨论】:

    标签: python apache-kafka kafka-consumer-api kafka-python


    【解决方案1】:

    一个topic的一个partition只能被同一个ConsumerGroup中的一个consumer消费。

    如果您不设置 group.id,KafkaConsumer 将为您生成一个新的随机 group.id。由于此 group.id 是唯一的,您将看到正在使用数据。

    如果您有多个使用相同 group.id 运行的消费者,则只有一个消费者会读取数据,而另一个则保持空闲状态,不消耗任何东西。

    【讨论】:

      【解决方案2】:

      我知道,这不是解决作者问题的方法。不过,如果你降落在这里,你可能会因为另一个原因遇到这个问题。和我一样。

      因此,至少对于 kafka-python v2.0.2 和 Aiven Kafka 代理设置,通过添加 consumer.poll() 的干调用解决了该问题。 这特别奇怪,因为在没有分配 group_id 时不需要这样做。

      输出来自:

      def get():
          for message in consumer:
              print(message.value)
          consumer.commit()
      

      在这种情况下什么都没有

      虽然下面按预期工作。它只从上次 commit() 中读出新消息:

      输出来自:

      def get():
          consumer.poll()
          for message in consumer:
              print(message.value)
          consumer.commit()
      

      它按预期输出自上次提交以来该主题中的所有消息

      JFYI,类构造函数如下所示:

          consumer = KafkaConsumer(
              topics,
              bootstrap_servers=self._service_uri,
              auto_offset_reset='earliest',
              enable_auto_commit=False,
              client_id='my_consumer_name',
              group_id=self.GROUP_ID,
              security_protocol="SSL",
              ssl_cafile=self._ca_path,
              ssl_certfile=self._cert_path,
              ssl_keyfile=self._key_path,
          )
      

      ¯\_(ツ)_/¯

      【讨论】:

        猜你喜欢
        • 2014-09-24
        • 2021-06-06
        • 2016-05-15
        • 1970-01-01
        • 2020-10-16
        • 2018-09-12
        • 2014-12-31
        • 2020-11-15
        • 2022-08-13
        相关资源
        最近更新 更多