โ† ใ‚นใ‚ญใƒซไธ€่ฆงใซๆˆปใ‚‹
eco2-team

event-driven

by eco2-team

๐ŸŒฑ ์ด์ฝ”์—์ฝ”(Ecoยฒ) BE

โญ 0๐Ÿด 0๐Ÿ“… 2026ๅนด1ๆœˆ25ๆ—ฅ

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)

AspectKafkaRabbitMQRedis (Ecoยฒ)
์šฉ๋„๋Œ€์šฉ๋Ÿ‰ ์ŠคํŠธ๋ฆฌ๋ฐTask QueueSSE ์‹ค์‹œ๊ฐ„ ์ „์†ก
Latency~10ms~5ms<1ms
Fan-outConsumer GroupExchangePub/Sub
์˜์†์„ฑํŒŒํ‹ฐ์…˜ ๋กœ๊ทธQueueStreams + KV
์šด์˜ ๋ณต์žก๋„๋†’์Œ (ZK/KRaft)์ค‘๊ฐ„๋‚ฎ์Œ

Ecoยฒ ์„ ํƒ ์ด์œ : SSE๋Š” ๋‹ค์ˆ˜ ํด๋ผ์ด์–ธํŠธ์— ์ €์ง€์—ฐ fan-out ํ•„์š” โ†’ Redis Pub/Sub ์ตœ์ 

Reference Files

ํ•ต์‹ฌ ์ปดํฌ๋„ŒํŠธ

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

ใ‚นใ‚ณใ‚ข

็ทๅˆใ‚นใ‚ณใ‚ข

60/100

ใƒชใƒใ‚ธใƒˆใƒชใฎๅ“่ณชๆŒ‡ๆจ™ใซๅŸบใฅใ่ฉ•ไพก

โœ“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

ใƒฌใƒ“ใƒฅใƒผ

๐Ÿ’ฌ

ใƒฌใƒ“ใƒฅใƒผๆฉŸ่ƒฝใฏ่ฟ‘ๆ—ฅๅ…ฌ้–‹ไบˆๅฎšใงใ™