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 eventIdA 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.