aboutsummaryrefslogtreecommitdiff
path: root/order-service/app/kafka_producer.py
diff options
context:
space:
mode:
authoralex <[email protected]>2026-07-16 15:02:05 +0200
committeralex <[email protected]>2026-07-16 15:02:05 +0200
commit4bc19be227f95ec06d33bca1b19bfa60afef99c6 (patch)
treebb8612ccca658474138ada62a9ea392f0444f4c3 /order-service/app/kafka_producer.py
parent8e12cb9bb66adc7adaf6d276133d6f235d294d23 (diff)
downloadorder-inventory-system-4bc19be227f95ec06d33bca1b19bfa60afef99c6.tar.xz
order-inventory-system-4bc19be227f95ec06d33bca1b19bfa60afef99c6.zip
order is now sent to kafka
Diffstat (limited to 'order-service/app/kafka_producer.py')
-rw-r--r--order-service/app/kafka_producer.py20
1 files changed, 20 insertions, 0 deletions
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()