diff options
| author | alex <[email protected]> | 2026-07-17 18:16:35 +0200 |
|---|---|---|
| committer | alex <[email protected]> | 2026-07-17 18:16:35 +0200 |
| commit | 8e796c9cfcd65f6225a6ae3ec2a4419f265b0ebc (patch) | |
| tree | 1bf90385144be2e9f350c31da553b1ac306ee60b /inventory-service/app/kafka_producer.py | |
| parent | 5f6918db393ed78dd8f038e9e5f1ea85f811dc52 (diff) | |
| download | order-inventory-system-8e796c9cfcd65f6225a6ae3ec2a4419f265b0ebc.tar.xz order-inventory-system-8e796c9cfcd65f6225a6ae3ec2a4419f265b0ebc.zip | |
Diffstat (limited to 'inventory-service/app/kafka_producer.py')
| -rw-r--r-- | inventory-service/app/kafka_producer.py | 33 |
1 files changed, 33 insertions, 0 deletions
diff --git a/inventory-service/app/kafka_producer.py b/inventory-service/app/kafka_producer.py new file mode 100644 index 0000000..ede6012 --- /dev/null +++ b/inventory-service/app/kafka_producer.py @@ -0,0 +1,33 @@ +import json +from confluent_kafka import Producer +from app.core.config import settings + +_conf = {"bootstrap.servers": settings.KAFKA_BOOTSTRAP_SERVERS} +_producer = Producer(_conf) + + +def _delivery_report(err, msg): + if err is not None: + print(f"[Kafka] Delivery failed: {err}") + else: + print(f"[Kafka] Message delivered → {msg.topic()} [partition {msg.partition()}]") + + +def send_inventory_event(event_payload: dict) -> None: + """ + Publish an inventory result event to the inventory topic. + Called after the consumer decides to emit StockReserved or StockReservationFailed. + + Args: + event_payload: A dict that conforms to either StockReservedEvent or + StockReservationFailedEvent schema. + """ + pass # TODO: implement + + +def flush_producer() -> None: + """ + Block until all buffered messages are delivered. + Call this on application shutdown. + """ + pass # TODO: implement |
