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