diff options
Diffstat (limited to 'order-service/app')
| -rw-r--r-- | order-service/app/api/v1/endpoints/orders.py | 8 | ||||
| -rw-r--r-- | order-service/app/init_db.py | 2 | ||||
| -rw-r--r-- | order-service/app/kafka_producer.py | 20 |
3 files changed, 28 insertions, 2 deletions
diff --git a/order-service/app/api/v1/endpoints/orders.py b/order-service/app/api/v1/endpoints/orders.py index 4530958..dc81835 100644 --- a/order-service/app/api/v1/endpoints/orders.py +++ b/order-service/app/api/v1/endpoints/orders.py @@ -5,6 +5,8 @@ from sqlalchemy.orm import Session from app.api import deps from app.models.order import Order, OrderItem from app.schemas.order import OrderCreate, OrderOut +from app.kafka_producer import send_order_created_event + router = APIRouter() @@ -38,10 +40,14 @@ def create_order(order_in: OrderCreate, db: Session = Depends(deps.get_db)): db.commit() db.refresh(db_order) + event_payload = {"order_id": db_order.id, "customer_id": order_in.customer_id} + send_order_created_event(event_payload) + return db_order - except Exception: + except Exception as e: db.rollback() + print(e) raise HTTPException( status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail="Could not initialize checkout process.", diff --git a/order-service/app/init_db.py b/order-service/app/init_db.py index 7ee1465..88a3149 100644 --- a/order-service/app/init_db.py +++ b/order-service/app/init_db.py @@ -1,6 +1,6 @@ from app.core.database import engine, Base -from app.models.order import Order, OrderItem +from app.models.order import Order, OrderItem # noqa: F401 def init_database(): diff --git a/order-service/app/kafka_producer.py b/order-service/app/kafka_producer.py new file mode 100644 index 0000000..6b24395 --- /dev/null +++ b/order-service/app/kafka_producer.py @@ -0,0 +1,20 @@ +import json +from confluent_kafka import Producer + +conf = {"bootstrap.servers": "localhost:9092"} +producer = Producer(conf) + + +def delivery_report(err, msg): + if err is not None: + print(f"Message delivery failed: {err}") + else: + print(f"Message delivered to {msg.topic()} [{msg.partition()}]") + + +def send_order_created_event(order_data: dict): + data_str = json.dumps(order_data) + + producer.produce("orders", value=data_str.encode("utf-8"), callback=delivery_report) + + producer.flush() |
