Skip to content

Commit 6ee08ac

Browse files
Create redis_bus.py
1 parent 9f84f78 commit 6ee08ac

1 file changed

Lines changed: 36 additions & 0 deletions

File tree

‎redis_bus.py‎

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,36 @@
1+
import json
2+
import asyncio
3+
import aioredis
4+
from .base import EventBus
5+
6+
7+
class RedisEventBus(EventBus):
8+
"""
9+
Redis Streams based implementation.
10+
Great for lightweight event-driven pipelines.
11+
"""
12+
13+
def __init__(self, redis_url="redis://localhost:6379"):
14+
self.redis_url = redis_url
15+
self.redis = None
16+
17+
async def connect(self):
18+
self.redis = await aioredis.from_url(self.redis_url, decode_responses=True)
19+
20+
async def publish(self, topic: str, message: dict):
21+
await self.redis.xadd(topic, {"data": json.dumps(message)})
22+
23+
async def subscribe(self, topic: str, handler):
24+
last_id = "$"
25+
while True:
26+
streams = await self.redis.xread({topic: last_id}, timeout=5000)
27+
if streams:
28+
_, messages = streams[0]
29+
for msg_id, fields in messages:
30+
last_id = msg_id
31+
data = json.loads(fields["data"])
32+
await handler(data)
33+
34+
async def close(self):
35+
if self.redis:
36+
await self.redis.close()

0 commit comments

Comments
 (0)