aboutsummaryrefslogtreecommitdiff
path: root/inventory-service/app/kafka_consumer.py
diff options
context:
space:
mode:
authoralex <[email protected]>2026-07-17 18:16:35 +0200
committeralex <[email protected]>2026-07-17 18:16:35 +0200
commit8e796c9cfcd65f6225a6ae3ec2a4419f265b0ebc (patch)
tree1bf90385144be2e9f350c31da553b1ac306ee60b /inventory-service/app/kafka_consumer.py
parent5f6918db393ed78dd8f038e9e5f1ea85f811dc52 (diff)
downloadorder-inventory-system-main.tar.xz
order-inventory-system-main.zip
inventory-service initHEADmain
Diffstat (limited to 'inventory-service/app/kafka_consumer.py')
-rw-r--r--inventory-service/app/kafka_consumer.py36
1 files changed, 36 insertions, 0 deletions
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()