โ Back to list

event-driven
by eco2-team
๐ฑ ์ด์ฝ์์ฝ(Ecoยฒ) BE
โญ 0๐ด 0๐
Jan 25, 2026
SKILL.md
name: event-driven description: Redis ๊ธฐ๋ฐ Composite Event Bus ํจํด ๊ฐ์ด๋. Event Router, SSE Gateway, Redis Streams/Pub-Sub ๊ตฌํ ์ ์ฐธ์กฐ. "event", "sse", "stream", "pubsub", "broadcast", "realtime" ํค์๋๋ก ํธ๋ฆฌ๊ฑฐ.
Event-Driven Architecture Guide
Ecoยฒ Composite Event Bus ์ํคํ ์ฒ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
โ Composite Event Bus Architecture โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโค
โ โ
โ Producers (Workers) โ
โ โโ scan-worker (Celery) โ Redis Streams (Durable) โ
โ โโ chat-worker (LangGraph) โ {domain}:events:{shard} โ
โ โโ character-worker โ
โ โ โ
โ โผ โ
โ โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ โ
โ โ Event Router (Event Bus Layer) โ โ
โ โ โโ Consumer: XREADGROUP (Consumer Group) โ โ
โ โ โโ Processor: Lua Script (State + Pub/Sub coordination) โ โ
โ โ โโ Reclaimer: XAUTOCLAIM (Fault Recovery) โ โ
โ โ โโ Publisher: Redis Pub/Sub (Real-time delivery) โ โ
โ โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ โ
โ โ โ โ
โ Streams Redis Pub/Sub Redis โ
โ (Durable Buffer) (Real-time Fan-out) โ
โ โโ {domain}:state:{job_id} โโ sse:events:{job_id} โ
โ โ โ โ
โ โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ โ
โ โ โ
โ โผ โ
โ โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ โ
โ โ SSE Gateway (Consumer/Publisher) โ โ
โ โ โโ Subscribe: Redis Pub/Sub channels โ โ
โ โ โโ Recovery: State KV polling โ โ
โ โ โโ Output: EventSource streaming โ โ
โ โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ โ
โ โ โ
โ โผ โ
โ Clients (Frontend) โ EventSource API โ
โ โ
โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
์ Redis์ธ๊ฐ? (vs Kafka/RabbitMQ)
| Aspect | Kafka | RabbitMQ | Redis (Ecoยฒ) |
|---|---|---|---|
| ์ฉ๋ | ๋์ฉ๋ ์คํธ๋ฆฌ๋ฐ | Task Queue | SSE ์ค์๊ฐ ์ ์ก |
| Latency | ~10ms | ~5ms | <1ms |
| Fan-out | Consumer Group | Exchange | Pub/Sub |
| ์์์ฑ | ํํฐ์ ๋ก๊ทธ | Queue | Streams + KV |
| ์ด์ ๋ณต์ก๋ | ๋์ (ZK/KRaft) | ์ค๊ฐ | ๋ฎ์ |
Ecoยฒ ์ ํ ์ด์ : SSE๋ ๋ค์ ํด๋ผ์ด์ธํธ์ ์ ์ง์ฐ fan-out ํ์ โ Redis Pub/Sub ์ต์
Reference Files
- Event Router: See event-router.md
- SSE Gateway: See sse-gateway.md
- Idempotency ํจํด: See idempotency.md
- Failure Recovery: See failure-recovery.md
- ACK Policy: See ack-policy.md โ ๏ธ Critical
- Reclaimer ํจํด: See reclaimer-patterns.md
- SSE ํ์ค: See sse-standard.md
ํต์ฌ ์ปดํฌ๋ํธ
1. Event Publishing (Worker โ Streams)
# Idempotent XADD with Lua Script
async def publish_stage_event(
job_id: str,
stage: str,
seq: int,
status: str,
data: dict,
) -> bool:
"""๋ฉฑ๋ฑ์ฑ ๋ณด์ฅ ์ด๋ฒคํธ ๋ฐํ"""
shard = md5_hash(job_id) % SHARD_COUNT
stream_key = f"{domain}:events:{shard}"
publish_key = f"published:{job_id}:{stage}:{seq}"
# Lua Script: Check โ XADD โ Mark
result = await redis.eval(
IDEMPOTENT_XADD_SCRIPT,
keys=[stream_key, publish_key],
args=[event_json, TTL],
)
return result == 1
2. Event Router Consumer
โ ๏ธ Critical: ACK๋ ์ฒ๋ฆฌ ์ฑ๊ณต ์์๋ง - See ack-policy.md
async def consume_loop():
"""XREADGROUP ๊ธฐ๋ฐ ์ด๋ฒคํธ ์๋น - ACK on Success Only"""
while True:
messages = await redis.xreadgroup(
groupname="eventrouter",
consumername=f"router-{POD_ID}",
streams={f"{domain}:events:{i}": ">" for i in range(SHARDS)},
block=5000,
count=100,
)
for stream, events in messages:
for event_id, data in events:
data["stream_id"] = event_id # SSE id ํ๋์ฉ
try:
success = await processor.process(data)
if not success:
continue # ACK ์คํต - PEL์ ์ ์ง
except Exception:
continue # ACK ์คํต - PEL์ ์ ์ง
await redis.xack(stream, "eventrouter", event_id) # ์ฑ๊ณต ์๋ง
3. SSE Gateway Streaming
SSE ํ์ค id: ํ๋๋ฅผ ํฌํจํ์ฌ Last-Event-ID ๋ณต๊ตฌ ์ง์ - See sse-standard.md
async def stream_events(job_id: str) -> AsyncGenerator[str, None]:
"""SSE ์ด๋ฒคํธ ์คํธ๋ฆฌ๋ฐ (ํ์ค id ํ๋ ํฌํจ)"""
channel = f"sse:events:{job_id}"
pubsub = redis.pubsub()
await pubsub.subscribe(channel)
# Recovery: State KV์์ ๋ง์ง๋ง ์ํ ์กฐํ
state = await redis.get(f"{domain}:state:{job_id}")
if state:
event = json.loads(state)
stream_id = event.get("stream_id", "")
stage = event.get("stage", "message")
yield f"event: {stage}\nid: {stream_id}\ndata: {state}\n\n"
# Real-time: Pub/Sub ๊ตฌ๋
async for message in pubsub.listen():
if message["type"] == "message":
event = json.loads(message["data"])
stream_id = event.get("stream_id", "")
stage = event.get("stage", "message")
yield f"event: {stage}\nid: {stream_id}\ndata: {message['data']}\n\n"
์ค์
Event Router
@dataclass
class EventRouterSettings:
redis_streams_url: str # XREADGROUP, State KV
redis_pubsub_url: str # PUBLISH
consumer_group: str = "eventrouter"
shard_count: int = 4
xread_block_ms: int = 5000
xread_count: int = 100
reclaim_min_idle_ms: int = 300000 # 5๋ถ
state_ttl: int = 3600 # 1์๊ฐ
published_ttl: int = 7200 # 2์๊ฐ
SSE Gateway
@dataclass
class SSEGatewaySettings:
redis_streams_url: str # State KV
redis_pubsub_url: str # Subscribe
sse_keepalive_interval: float = 15.0
sse_max_wait_seconds: int = 300
sse_queue_maxsize: int = 100
Score
Total Score
60/100
Based on repository quality metrics
โSKILL.md
SKILL.mdใใกใคใซใๅซใพใใฆใใ
+20
โLICENSE
ใฉใคใปใณในใ่จญๅฎใใใฆใใ
+10
โ่ชฌๆๆ
100ๆๅญไปฅไธใฎ่ชฌๆใใใ
0/10
โไบบๆฐ
GitHub Stars 100ไปฅไธ
0/15
โๆ่ฟใฎๆดปๅ
3ใถๆไปฅๅ ใซๆดๆฐใใใ
0/10
โใใฉใผใฏ
10ๅไปฅไธใใฉใผใฏใใใฆใใ
0/5
โIssue็ฎก็
ใชใผใใณIssueใ50ๆชๆบ
+5
โ่จ่ช
ใใญใฐใฉใใณใฐ่จ่ชใ่จญๅฎใใใฆใใ
+5
โใฟใฐ
1ใคไปฅไธใฎใฟใฐใ่จญๅฎใใใฆใใ
0/5
Reviews
๐ฌ
Reviews coming soon