Repository navigation
Expand file tree
/
Copy pathkafka_bus.py
More file actions
49 lines (39 loc) · 1.5 KB
/
Copy pathkafka_bus.py
File metadata and controls
49 lines (39 loc) · 1.5 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
import json
from aiokafka import AIOKafkaConsumer, AIOKafkaProducer
from .base import EventBus
class KafkaEventBus(EventBus):
"""
Kafka implementation of the EventBus interface.
Uses aiokafka for async publish/subscribe.
"""
def __init__(self, bootstrap_servers="localhost:9092", group_id="ride-sharing"):
self.bootstrap_servers = bootstrap_servers
self.group_id = group_id
self.producer = None
self.consumer = None
async def connect(self):
self.producer = AIOKafkaProducer(
bootstrap_servers=self.bootstrap_servers,
value_serializer=lambda v: json.dumps(v).encode("utf-8"),
)
await self.producer.start()
async def publish(self, topic: str, message: dict):
if not self.producer:
raise RuntimeError("Kafka producer not initialized. Call connect().")
await self.producer.send_and_wait(topic, message)
async def subscribe(self, topic: str, handler):
consumer = AIOKafkaConsumer(
topic,
bootstrap_servers=self.bootstrap_servers,
group_id=self.group_id,
value_deserializer=lambda v: json.loads(v.decode("utf-8")),
)
await consumer.start()
try:
async for msg in consumer:
await handler(msg.value)
finally:
await consumer.stop()
async def close(self):
if self.producer:
await self.producer.stop()