From 4bc19be227f95ec06d33bca1b19bfa60afef99c6 Mon Sep 17 00:00:00 2001 From: alex Date: Thu, 16 Jul 2026 15:02:05 +0200 Subject: order is now sent to kafka --- order-service/app/api/v1/endpoints/orders.py | 8 +++++++- order-service/app/init_db.py | 2 +- order-service/app/kafka_producer.py | 20 ++++++++++++++++++++ order-service/pyproject.toml | 1 + order-service/uv.lock | 15 +++++++++++++++ 5 files changed, 44 insertions(+), 2 deletions(-) create mode 100644 order-service/app/kafka_producer.py 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() diff --git a/order-service/pyproject.toml b/order-service/pyproject.toml index 8828bd3..0c51696 100644 --- a/order-service/pyproject.toml +++ b/order-service/pyproject.toml @@ -3,6 +3,7 @@ name = "order-service" version = "0.1.0" requires-python = ">=3.14" dependencies = [ + "confluent-kafka>=2.15.0", "fastapi>=0.139.0", "psycopg2>=2.9.12", "pydantic>=2.13.4", diff --git a/order-service/uv.lock b/order-service/uv.lock index bc3cc5c..833c0ec 100644 --- a/order-service/uv.lock +++ b/order-service/uv.lock @@ -53,6 +53,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/d1/d6/3965ed04c63042e047cb6a3e6ed1a63a35087b6a609aa3a15ed8ac56c221/colorama-0.4.6-py2.py3-none-any.whl", hash = "sha256:4f1d9991f5acc0ca119f9d443620b77f9d6b33703e51011c16baf57afb285fc6", size = 25335, upload-time = "2022-10-25T02:36:20.889Z" }, ] +[[package]] +name = "confluent-kafka" +version = "2.15.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/51/90/eb998fedefb63b42910b54b76b7d300ccddc56430a5175122eb60cedc4f3/confluent_kafka-2.15.0.tar.gz", hash = "sha256:7ad9bad1cbabf6713ec039b8204b48d322024fd11397eec88d912e048c732ba7", size = 323365, upload-time = "2026-06-30T19:48:43.395Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/18/20/5c078b9560b8005404e09a5889b9771e8deb5c9dddeb9003f58c7d028258/confluent_kafka-2.15.0-cp314-cp314-macosx_13_0_arm64.whl", hash = "sha256:e781d39ff31d5d10d0502471f92f95aff4f5be1eaae9a1f57aeac164bc9f9029", size = 4343440, upload-time = "2026-06-30T19:48:17.324Z" }, + { url = "https://files.pythonhosted.org/packages/aa/27/fd83281231b6179e3923822e20861dfc1d6c200edb910c3d8e9a15b9e95a/confluent_kafka-2.15.0-cp314-cp314-macosx_13_0_x86_64.whl", hash = "sha256:000463c822a7adc3293369b63688b68d17c03a1a4a64f86182efc08f8b150676", size = 4371051, upload-time = "2026-06-30T19:48:19.29Z" }, + { url = "https://files.pythonhosted.org/packages/5e/3e/fa82e6699144e707972bbf57489598defdf67e844ae2a7d29e0ea4e3a187/confluent_kafka-2.15.0-cp314-cp314-manylinux_2_28_aarch64.whl", hash = "sha256:9ddf4cf4647e5d633ef64e3f5a349d5288edfedf974203993e9d86548c2695be", size = 4979624, upload-time = "2026-06-30T19:48:21.454Z" }, + { url = "https://files.pythonhosted.org/packages/13/ee/b4b6a0da17584b432c83a0500ac79c04864ef92dafdb405614831b995232/confluent_kafka-2.15.0-cp314-cp314-manylinux_2_28_x86_64.whl", hash = "sha256:d8ed33f623ff2a104fb76f99a67e9917f0170fddb4e28380fcdc83347b1646b2", size = 4785678, upload-time = "2026-06-30T19:48:23.068Z" }, + { url = "https://files.pythonhosted.org/packages/a8/f4/19a852d16e8e4f8ac930037d8ecda21a220f6d5d050a59bbce10189ac9ec/confluent_kafka-2.15.0-cp314-cp314-win_amd64.whl", hash = "sha256:2e80bd96f61aae2ffba951754a769e6d3c5ebb5a5e778ed0ae8ad899ea91556b", size = 4744444, upload-time = "2026-06-30T19:48:24.698Z" }, +] + [[package]] name = "fastapi" version = "0.139.0" @@ -131,6 +144,7 @@ name = "order-service" version = "0.1.0" source = { virtual = "." } dependencies = [ + { name = "confluent-kafka" }, { name = "fastapi" }, { name = "psycopg2" }, { name = "pydantic" }, @@ -141,6 +155,7 @@ dependencies = [ [package.metadata] requires-dist = [ + { name = "confluent-kafka", specifier = ">=2.15.0" }, { name = "fastapi", specifier = ">=0.139.0" }, { name = "psycopg2", specifier = ">=2.9.12" }, { name = "pydantic", specifier = ">=2.13.4" }, -- cgit v1.2.3