aboutsummaryrefslogtreecommitdiff
path: root/order-service/app/kafka_producer.py
blob: 6b24395916b1583566d3d99f0d30483460b4f51c (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
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()