skip to content

When should an in-process Observer be replaced by publish-subscribe through a message broker, and what changes about the guarantees?

level: principalimportance: should knowfreq 40%

answer

  1. Observer decouples types; broker also decouples time + lifetime
  2. in-process: sync, exactly-once, dies with the process
  3. broker: durable, replayable, at-least-once ⇒ idempotency
  4. ordering only per partition/key; dual write ⇒ outbox
  5. costs: schema contract, dead letters, lag, eventual consistency

basics

~20 s

Use in-process Observer when subject and observers live in the same process and must react immediately. Move to a broker when the reactors are separate services or must survive restarts. You gain durability, retries and independent scaling; you lose synchronous ordering and immediate consistency, and you must handle duplicates.

solid answer

~60 s

In-process Observer decouples types but not time, lifetime or failure: notification is a synchronous call on the caller's thread, observers are references the subject holds, and everything dies with the process. A broker adds an intermediary that decouples all three — publisher and subscriber need not be running at the same time, know each other's addresses, or share a process. The trade is a different set of guarantees. You get durability, retry, dead-lettering, independent scaling and consumption by teams you have never met. You give up: exactly-once (you get at-least-once, so consumers must be idempotent, usually via a dedup key), global ordering (typically ordered only per partition/key), read-your-writes (consumers are eventually consistent), and the simple stack trace — failures now appear as lag and dead-letter depth. You also inherit the dual-write problem: publishing after a commit can lose events, publishing inside the transaction can emit events for rolled-back state, so you need a transactional outbox or equivalent. Choose in-process for intra-aggregate/UI consistency; choose a broker for cross-service, cross-team, or replayable integration.

code

pseudocode · 16 lines
pseudocode
// In-process: immediate, atomic with the change, dies with the process
order.place()
  -> repository.save(order)
  -> notify(OrderPlaced(order.id))     // same thread, same transaction

// Cross-service: transactional outbox avoids the dual-write problem
transaction {
    repository.save(order)
    outbox.insert(Event(id = uuid(), key = order.id, type = "OrderPlaced", version = 1, payload))
}   // state + intent commit atomically

relay: poll outbox -> publish to topic (partition key = order.id) -> mark sent   // at-least-once

consumer.handle(e):
    if (processed.contains(e.id)) return          // idempotency by event id
    apply(e); processed.add(e.id)                 // in one transaction where possible

go deeper

for a junior

Say that Observer works within one program while a message broker lets separate services react, and that broker delivery is asynchronous and can retry.

for a middle

Contrast the guarantees: in-process is synchronous, ordered and exactly-once but dies with the process; a broker is durable and asynchronous but at-least-once, so handlers must be idempotent.

for a senior

Add the engineering obligations: idempotency keys, per-partition ordering, the dual-write problem and transactional outbox, schema versioning, retries and dead-letter policy, consumer lag observability, and eventual consistency in the UI.

for a principal

Present it as an architectural boundary decision with organizational consequences: which reactions may cross a team boundary, who owns the event schema and its compatibility policy, what the failure and on-call model becomes, and when an in-process event bus is the correct intermediate step. Name the failure mode of getting it wrong in each direction — a distributed monolith versus needless infrastructure for a local method call.

## What actually differs Both are "one producer, many consumers", but they decouple different dimensions: | Dimension | In-process Observer | Broker pub/sub | |---|---|---| | Type coupling | Removed (subject sees an interface) | Removed (topic + schema) | | Temporal coupling | **Present** — synchronous call, publisher blocks | Removed — consumer may be down | | Lifetime coupling | **Present** — subject strongly references observers | Removed — subscriptions live in the broker | | Location | Same process/address space | Any process, host, team | | Ordering | Deterministic within one dispatch | Per-partition/key at best | | Delivery | Exactly once (it's a method call) | At-least-once, typically | | Durability | None — dies with the process | Persisted, replayable | | Failure surface | Exception in the caller's stack | Lag, retries, dead-letter queue | | Consistency | Immediate, possibly same transaction | Eventual | | Observability | Debugger/stack trace | Traces, correlation ids, consumer metrics | | Backpressure | None — publisher pays inline | Broker buffers; consumers lag | ## When in-process Observer is the right answer - Subject and reactors are in the same deployable and always co-live (a model and its views, a cache and its invalidator, a domain aggregate and its intra-module reaction). - The reaction must be visible immediately after the change returns (read-your-writes; a UI must not show stale state for a second). - The reaction is cheap and local (recompute a derived field, mark dirty, repaint). - You want debuggability: one stack trace explains everything. - Adding infrastructure would exceed the value: a broker brings deployment, schema management, monitoring, dead-letter policy, and on-call load. ## When to move to a broker - **Different deployables / teams.** A reaction owned by another service must not be a compile-time or runtime dependency of the writer. - **The publisher must not pay the cost.** Slow, bursty or unreliable reactions (indexing, exports, notifications, ML scoring) should not add latency to the write path or fail it. - **The reaction must survive restarts and outages.** Consumers can be down for an hour and catch up; in-process observers simply lose the events. - **Unknown or growing consumer set.** Broadcast to teams you will never coordinate with, with independent onboarding and replay from an offset. - **Replay / audit / rebuild.** A durable log lets you rebuild a projection or backfill a new consumer from history — impossible with in-memory notification. - **Independent scaling and backpressure.** Consumers scale by partition; a slow consumer accrues lag instead of blocking the producer. ## What you must now engineer (the honest cost list) 1. **Idempotency.** At-least-once delivery means duplicates on retry or rebalance. Every consumer needs a natural dedup key (event id, aggregate id + version) and idempotent side effects. "Exactly-once" claims are, in practice, at-least-once plus deduplication at the effect boundary. 2. **Ordering.** Global ordering does not exist at scale. Partition by an entity key so per-entity order holds, and design handlers to tolerate out-of-order or reordered delivery (version checks, last-write-wins with timestamps/versions). 3. **The dual-write problem.** "Commit to the database, then publish" can crash between the two (event lost); "publish, then commit" can emit events for state that rolls back (phantom event). The standard fix is the **transactional outbox**: write the event to a table in the same transaction, and a separate relay publishes it — restoring atomicity between state and event at the cost of eventual publication and at-least-once semantics. Change-data-capture off the write-ahead log is the same idea with less application code. 4. **Schema contract and evolution.** The event payload is now a published API consumed by code you cannot refactor. You need a schema registry or equivalent, backward/forward compatibility rules, and versioning discipline; adding a required field becomes a breaking change. 5. **Poison messages and dead letters.** Define retry policy with backoff, a maximum attempt count, a dead-letter destination, and — importantly — an operational procedure for draining it. A dead-letter queue nobody reads is a silent data-loss channel. 6. **Observability.** Replace the stack trace with correlation/trace ids propagated through events, consumer lag dashboards, dead-letter depth alarms, and end-to-end latency SLOs. 7. **Eventual consistency in the product.** The UI must tolerate "submitted, not yet reflected": show pending states, avoid read-after-write assumptions, or read from the write side for the user's own actions. 8. **Semantic care in the event itself.** Prefer facts about the past (`OrderPlaced`) over commands disguised as events (`SendEmail`); include the aggregate identity and version; decide between thin events (id only, consumer fetches — risks reading newer state) and fat events (self-contained snapshot — bigger, but consistent and replayable). ## The middle ground worth mentioning An **in-process event bus/mediator** (still Observer, but with an intermediary) buys you indirection and cross-cutting behavior — async hand-off to a bounded executor, per-handler error isolation, metrics, ordering policy — without introducing infrastructure. It is very often the right next step before a broker, and it makes the eventual migration mechanical because publishers already emit event values instead of calling listeners. ## The decision rule to state out loud *Keep it in-process while the reaction must be immediate, local, and atomic with the change; move to a broker when the reaction must outlive the process, cross a team boundary, or scale independently — and accept that this converts a method call into a distributed-systems problem with idempotency, ordering, schema and operational obligations.* Making that trade for a reaction that could have been three lines in the same process is over-engineering; refusing it for a cross-team integration produces a distributed monolith where one team's deploy breaks another's write path.

  • What is the dual-write problem and how does a transactional outbox solve it?
    Writing state to the database and publishing an event to a broker are two separate systems with no shared transaction, so a crash between them either loses the event or emits an event for state that rolled back. The outbox writes the event into a table inside the same database transaction as the state change, and a separate relay reads that table and publishes. Atomicity is restored between state and intent; publication becomes eventual and at-least-once, so consumers must be idempotent.
  • Why can't a broker give you true exactly-once delivery, and what do you do instead?
    Acknowledgements can be lost and consumers can crash after acting but before committing an offset, so any system that retries must be prepared to deliver again. Practical exactly-once is at-least-once delivery plus deduplication at the effect boundary: a stable event id or an (aggregate, version) pair, an idempotency table or upsert, and side effects designed so repeating them is harmless.
  • How do you decide between a fat event carrying the full state and a thin event carrying just an identifier?
    Fat events are self-contained, replayable, and give every consumer the same consistent snapshot, at the cost of payload size and a schema more consumers depend on. Thin events are small and keep the payload contract minimal, but every consumer must call back to the source — adding load, coupling to its query API, and the risk of reading a state newer than the event described. It is the same push-versus-pull trade-off as inside a single process, with network cost and staleness amplified.

Telling colleagues in the room versus posting to a company-wide mailing list with an archive. The room is instant, everyone hears the same thing once, and nothing is recorded; the mailing list reaches people who are on holiday, keeps a searchable history that newcomers can read from the beginning — and occasionally delivers the same message twice, out of order, to someone who has already acted on it.

saying these in an interview costs you the question

  • Claiming a broker gives exactly-once delivery and global ordering.
  • Introducing a broker for reactions that live in the same process and must be immediately consistent — infrastructure cost with no decoupling gained.
  • Publishing an event before the state change commits, or after it, without an outbox — the dual-write problem.
  • Treating the event payload as an internal detail rather than a versioned public contract.
  • Assuming in-process Observer gives temporal decoupling; it is a synchronous call that blocks the publisher.
  • No dead-letter policy, or a dead-letter queue nobody monitors or drains.
  • Publishing commands ('SendEmail') as events instead of facts about what happened, which recreates coupling through the topic.

context