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()
|