aboutsummaryrefslogtreecommitdiff
path: root/inventory-service/app/kafka_consumer.py
blob: 9fda00db47af77c3325e75a12fa840f96b8a8e20 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
from confluent_kafka import Consumer, KafkaException
from app.core.config import settings
from app.core.database import SessionLocal


def _build_consumer() -> Consumer:
    conf = {
        "bootstrap.servers": settings.KAFKA_BOOTSTRAP_SERVERS,
        "group.id": settings.KAFKA_CONSUMER_GROUP_ID,
        "auto.offset.reset": "earliest",
        "enable.auto.commit": False,
    }
    consumer = Consumer(conf)
    consumer.subscribe([settings.KAFKA_ORDERS_TOPIC])
    return consumer


def handle_order_created(message_value: dict, db) -> None:

    pass  # TODO: implement


def start_consumer_loop() -> None:

    consumer = _build_consumer()
    db = SessionLocal()

    try:
        while True:
            # TODO: poll, decode, call handle_order_created, commit offset
            pass
    except KafkaException as e:
        print(f"[Kafka Consumer] Fatal error: {e}")
    finally:
        consumer.close()
        db.close()