Lesson 17 / 25

Event Streams Between Services

Use streams to publish domain events with an outbox and idempotent consumers.

A lightweight event bus

Streams work well as a lightweight event bus between a modest number of services. A service appends domain events (OrderPlaced, PaymentCaptured) to a stream per domain, such as events:orders; each interested service creates its own consumer group and processes events at its own pace, with acknowledgements and recovery. New services can start a group from 0 to process the retained history. Reliability rules from event-driven design still apply: write events through a transactional outbox so database changes and published events cannot disagree (a relay reads the outbox table and calls XADD); include an event ID and type in each entry; make consumers idempotent, since redelivery happens; and keep retention long enough for consumers to recover from outages. Because a stream is one key, all events in it are ordered and live on a single shard; partition across several streams when one becomes too hot.

An outbox relay publishing to a stream

The relay marks rows as published only after XADD succeeds.

def relay_outbox(db, r, batch=100):
    rows = db.fetch_all(
        "SELECT id, aggregate_id, type, payload FROM outbox "
        "WHERE published_at IS NULL ORDER BY id LIMIT %s FOR UPDATE SKIP LOCKED", (batch,))
    for row in rows:
        r.xadd("events:orders",
               {"eventId": str(row.id), "type": row.type,
                "aggregateId": row.aggregate_id, "payload": row.payload},
               maxlen=1_000_000, approximate=True)
        db.execute("UPDATE outbox SET published_at = now() WHERE id = %s", (row.id,))
    db.commit()
# a crash between XADD and UPDATE republishes the event; consumers dedupe by eventId

A shared logbook in a control room

Each department writes significant events in the control-room logbook. Every other department keeps its own bookmark and reads new lines when ready; nobody tears pages out, and the oldest pages are archived after a while.

Quick check: How should several services each receive every event from the same Redis stream?

  • Use one consumer group shared by all services
  • Give each service its own consumer group
  • Use XDEL after each read
  • Publish a copy of each event per service
Answer

Give each service its own consumer group — Each group tracks its own position, so every service sees every event.