Я использую confluent-kafka-client. У меня есть один производитель, создающий тему с одним разделом и одним потребителем в одном идентификаторе группы. Сначала я создаю производителя (с конфигурациями по умолчанию) для темы (если темы не существует, я создаю ее с таким именем)
Код: Выделить всё
self.producer = confluent_kafka.Producer({"bootstrap.servers": bootstrap_servers})
Затем я создаю потребителя и подписываю его на тему (с конфигурациями по умолчанию, auto.offset.reset="latest")
Код: Выделить всё
self.consumer = confluent_kafka.Consumer(
{"bootstrap.servers": self.bootstrap_servers,
"group.id": self.group_id},
logger=logger,
)
self.consumer.subscribe(self.topic_names, on_assign=print_assignment)
self.consumer.poll(0) # first call
Я понял, что self.consumer.poll(0) не регистрирует этого потребителя в теме, поскольку данных по этой теме еще нет. После этого продюсер выпускает пластинку. Затем я звоню
ожидаем получения данных. Однако он возвращает None. Фактически, после создания данных вызов poll(0) регистрирует потребителя. Я могу получить данные, позвонив в третий раз.
Как зарегистрировать потребителя в теме, если по этой теме еще нет данных?
Ссылки: возвращает Kafka Consumer.poll нет записей
Подробнее здесь:
https://stackoverflow.com/questions/783 ... ns-no-data