Skip to main content

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)

Топик в списке появляется не сразу, а где-то через 30 секунд. 

Список тем: 

def list_topics():
    metadata = admin_client.list_topics(timeout=10)
    print("Существующие топики:", list(metadata.topics.keys()))

consumer.py

from confluent_kafka import Consumer, KafkaException

config = {
    'bootstrap.servers': '192.168.1.195:29092',
    'group.id' : 'my_group'
}

consumer = Consumer(config)

if __name__ == '__main__':
    topic = 'test_topic'
    timeout_s = 1
    consumer.subscribe([topic])

    try:
        while True:
            msg = consumer.poll(timeout=timeout_s)
            if msg is None:
                continue
            if msg.error():
                raise KafkaException(msg.error())
            else:
                print(f'Input: {msg.value().decode()}')
    finally:
        consumer.close()

Если consumer в одной группе и количество partitions не меньше количества consumers, то сообщения распределяются между ними. Если partitions будет меньше consumers (которые в одной группе), то оставшаяся часть будет простаивать.

Функция consumer.subscribe получает список топиков, то есть одним consumer можно подписаться на несколько.

Если при двух и более partitions для одного из partitions отваливается consumer то Kafka автоматически проводит перебалансировку в пределах группы. После подключения опять произойдет перебалансировка.

producer.py

import time
from confluent_kafka import Producer

def delivery_report(err, msg):
    if err is not None:
        print(f"Message failed: {err}")
        raise err
    else:
        print(f'Message {msg.value()} in topic {msg.topic()} delivered')

config = {
    'bootstrap.servers': '192.168.1.195:29092',
    'message.timeout.ms' : 5000
}

producer = Producer(config)

def send_message(topic_: str, message_: str):
    producer.produce(topic_, value=message_, callback=delivery_report)
    producer.flush()

if __name__ == '__main__':
    topic = 'test_topic'
    message = 'Producer 1'
    
    i = 0
    while True:
        send_message(topic, f'{message} # {i}')
        i += 1
        time.sleep(1)

Распределение сообщений по ключам

Для распределения по консюмерам используется хэш функция от ключа и количества партиций. При коротких ключах может быть коллизия. Довольно неприятная и сложнообходимая ситуация. Проще указать номер partition (int).

import time
from confluent_kafka import Producer

def delivery_report(err, msg):
    if err is not None:
        print(f"Message failed: {err}")
        raise err
    else:
        print(f'Message {msg.value()} in topic {msg.topic()} delivered')

config = {
    'bootstrap.servers': '192.168.1.195:29092',
    'message.timeout.ms' : 5000
}

producer = Producer(config)

def send_message(topic_: str, message_: str, part_: int | None = None):
    producer.produce(topic_, partition=part_, value=message_, callback=delivery_report)
    producer.flush()

if __name__ == '__main__':
    topic = 'test_topic'
    message = 'Producer 1'
    
    i = 0
    while True:
        part1 = 0
        part2 = 1
        send_message(topic, message_=f'key={key1} {message} # {i}', part_=part1)
        send_message(topic, message_=f'key={key2} {message} # {i}', part_=part2)
        i += 1
        time.sleep(1)

#============================================ consumer =========================================
from confluent_kafka import Consumer, KafkaException

config = {
    'bootstrap.servers': '192.168.1.195:29092',
    'group.id' : 'my_group'
}

consumer = Consumer(config)

if __name__ == '__main__':
    topic = 'test_topic'
    timeout_s = 1
    consumer.subscribe([topic])

    try:
        while True:
            msg = consumer.poll(timeout=timeout_s)
            if msg is None:
                continue
            if msg.error():
                raise KafkaException(msg.error())
            else:
                print(f'Input: {msg.value().decode()}')
    finally:
        consumer.close()

Транзакции

В config для consumer добавляется элемент 'enable.auto.commit': False

И затем consumer.commit()

Сериализация данных

Через json 

while True:
        obj = {
            'id': i,
            'value': randomword(8),
            'timestamp': datetime.now().isoformat(timespec='seconds'),
        }

        send_message(topic, json.dumps(obj))
        i += 1
        time.sleep(1)
============================================================================
                str_msg = msg.value().decode('utf-8')
                obj = json.loads(str_msg)
                print(f'id={obj.get("id")}, timestamp={obj.get("timestamp")}')

 

 

 

 

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())