Python + Kafka
Библиотеки
- confluent_kafka от разработчиков
- aiokafka асинхронная библиотека
confluent_kafka
AdminClient
Модуль автоматизации административных задач.
Настройки:
from confluent_kafka.admin import AdminClient, NewTopic
config = {
'bootstrap.servers': '192.168.1.195:29092'
}
admin_client=AdminClient(config)
Создание темы:
def create_partition():
topic = 'test_topic'
new_topic = NewTopic(topic, num_partitions=2, replication_factor=1)
futures = admin_client.create_topics([new_topic])
for topic_name, future in futures.items():
try:
future.result()
print(f"Created {topic_name} successfully.")
except Exception as e :
print(f'Error in topic {topic_name} creation: ', e)
Список тем:
def list_topics():
metadata = admin_client.list_topics(timeout=10)
print("Существующие топики:", list(metadata.topics.keys()))
consumer.py
producer.py
aiokafka
consumer.py
from aiokafka import AIOKafkaConsumer
import asyncio
async def consume():
consumer = AIOKafkaConsumer(
'test_first_one',
bootstrap_servers='192.168.1.195:29092'
)
await consumer.start()
try:
async for msg in consumer:
print(f"Consumed message {msg.value.decode()}")
finally:
await consumer.stop()
if __name__ == "__main__":
asyncio.run(consume())
producer.py
from aiokafka import AIOKafkaProducer
import asyncio
async def send():
producer = AIOKafkaProducer(
bootstrap_servers="192.168.1.195:29092"
)
topik = 'test_first_one'
message = 'Hello from python!'
await producer.start()
i = 0
try:
while True:
bytes_msg = f'{message} {i}'.encode('utf-8') # convert string into byte
await producer.send_and_wait(topik, bytes_msg)
i += 1
print(f'Message {i} sent.')
await asyncio.sleep(1)
finally:
await producer.stop()
if __name__ == '__main__':
asyncio.run(send())