From 8e796c9cfcd65f6225a6ae3ec2a4419f265b0ebc Mon Sep 17 00:00:00 2001 From: alex Date: Fri, 17 Jul 2026 18:16:35 +0200 Subject: inventory-service init --- inventory-service/app/kafka_consumer.py | 36 +++++++++++++++++++++++++++++++++ 1 file changed, 36 insertions(+) create mode 100644 inventory-service/app/kafka_consumer.py (limited to 'inventory-service/app/kafka_consumer.py') diff --git a/inventory-service/app/kafka_consumer.py b/inventory-service/app/kafka_consumer.py new file mode 100644 index 0000000..9fda00d --- /dev/null +++ b/inventory-service/app/kafka_consumer.py @@ -0,0 +1,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() -- cgit v1.2.3