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