aio-libs / aio-libs/aiokafka

[QUESTION] Which Partition to use when Performing a Offset Seek using AIOKafkaConsumer?

Aberta
#772 6 comentários 0 reações 0 responsáveis Ver no GitHub
question
Linguagem predominante
Python
Estrelas
1.4k
Forks
269
Merge médio
1d 1h
PRs com merge (30d)
6

Descrição

When trying to let an `AIOKafkaConsumer` start reading messages from a specific offset `starting_offset`, how do we know which partition to be used?

I am trying to use the [`AIOKafkaConsumer.seek`][1] method, but it requires a _TopicPartition_ to be specified in.

```py
import asyncio
from aiokafka import AIOKafkaConsumer, AIOKafkaProducer

async def main():
topic = "test"
starting_offset = 3

# Publish some messages
producer = AIOKafkaProducer(bootstrap_servers="localhost:29092")
await producer.start()
for i in range(10):
await producer.send_and_wait(topic, bytes(f"hello {i}", "utf-8"))

# Start consuming from a specific offset
consumer = AIOKafkaConsumer(topic, bootstrap_servers="localhost:29092")
await consumer.start()
consumer.seek(None, starting_offset) # NEED HELP HERE :)

while True:
message = await consumer.getone()
print("message:", message.value)

if __name__ == "__main__":
asyncio.run(main())
```

Thanks!

[1]: https://aiokafka.readthedocs.io/en/stable/api.html#aiokafka.AIOKafkaConsumer.seek

Guia de contribuição

Abrir o guia de contribuição

Avaliação

Esta issue ainda não foi avaliada.

Receba novas issues na sua caixa de entrada

Um resumo curto de issues do GitHub para quem está começando.