Publish-subscribe and event notification
Chapter 7 introduced OrderPlaced as a Kafka event consumed independently by inventory-service and payment-service. This chapter makes explicit a design decision that chapter glossed over: exactly how much data belongs inside that event, and why the answer differs by consumer.
1. Problem the Pattern Solves
Two months after Chapter 7 ships, the newly added loyalty-service (from that chapter’s exercise) needs to know the customer’s tier (GOLD, SILVER, STANDARD) to calculate points correctly — a field OrderPlaced doesn’t carry. The loyalty team’s fastest fix is to add a synchronous REST call back to order-service to fetch the customer’s tier whenever an OrderPlaced event arrives — quietly reintroducing exactly the temporal coupling Chapter 7 was built to remove, just one hop later than before.
Meanwhile, a different consumer, analytics-pipeline, wants the entire order payload — every line item, every price, shipping address, the works — to avoid a callback of its own. If order-service puts all of that directly into OrderPlaced to satisfy analytics-pipeline, inventory-service and payment-service now receive a payload ten times larger than they need, and every unrelated field order-service ever adds to an order becomes a change to a payload three other services deserialize.
Forces in tension:
- Coupling to the publisher’s schema vs. self-sufficiency. A large event lets consumers avoid ever calling back to the publisher, but tightly couples every consumer to the exact shape of the publisher’s internal data, which now can’t change without considering every subscriber.
- Payload size and broker load vs. round-trip elimination. Small events keep the broker’s throughput high and payloads cheap to serialize, but may force a consumer to make a follow-up call for data it needs, reintroducing a dependency the event was meant to avoid.
- Consumer autonomy vs. staleness risk. A “fat” event lets a consumer build its own local copy of exactly the data it needs (as
loyalty-service’s tier lookup, done right, should) — but that local copy can drift from the source of truth if update events are ever missed or misordered.
2. Core Idea
Event notification: a thin event carrying only an identifier and the fact that something happened (OrderPlaced(orderId: UUID)), leaving it to interested consumers to call back to the publisher’s API if they need more detail. Simple, small, and minimizes schema coupling — at the cost of reintroducing a synchronous dependency for any consumer that needs data beyond the identifier.
Event-carried state transfer: a “fatter” event carrying enough data for consumers to act without calling back (OrderPlaced(orderId, items, customerId, customerTier, totalAmount)), letting each consumer maintain its own local, denormalized copy of the data it needs. Removes the callback dependency entirely — at the cost of a wider, more consequential schema that more consumers depend on directly.
Both are publish-subscribe — any number of independent subscribers, no publisher awareness of who’s listening (unchanged from Chapter 7) — the difference is purely in payload design, and Northwind, correctly, uses different answers for different consumers of the same underlying fact.
Commonly confused with:
- Event sourcing (a later chapter). Event-carried state transfer means a consumer builds a denormalized read model from events it receives; event sourcing means a service’s own authoritative state is the event log, replayable from scratch.
loyalty-service’s local tier cache in this chapter is a read model built from events — not event sourcing, sinceorder-service(the source of truth for orders) still stores current-state rows, not an event log, as established in Chapter 1. - CQRS (a later chapter). A consumer maintaining its own read model from events is a building block CQRS often uses, but CQRS specifically means separating a single service’s own read and write models — this chapter is about data flowing between services, not within one.
- A shared database view. Building a local denormalized copy from events looks superficially similar to querying another service’s database directly, but the event is a versioned, intentional, publisher-controlled contract — a shared view is an accidental coupling to internal storage structure that Chapter 2 already ruled out.
3. When to Use It
Use thin event notification when:
- Few consumers need the extra data, and they need it rarely enough that an occasional callback is cheap relative to bloating every event for every consumer.
- The data referenced changes frequently or is large (e.g., a full customer profile) — carrying it in every event risks staleness or payload bloat for consumers who rarely need it.
- Strong consistency for that specific field matters more than avoiding a callback — a callback gets the current value; an embedded value in an old event might be stale.
Use event-carried state transfer when:
- A consumer’s core job depends on this data for every event, at high volume —
loyalty-servicecalculating points for every single order is exactly this case; a callback per order recreates the coupling this pattern exists to avoid. - The publisher can tolerate consumers holding slightly stale, eventually-consistent copies of the embedded data (a customer’s tier at the moment of ordering, not necessarily their tier right now).
- You want to eliminate a specific, identified synchronous dependency that has already caused (or is about to cause) an availability problem — not speculatively, for every event.
Concrete use cases:
- E-commerce: shipping notifications carrying enough address and item detail for a carrier-integration service to act without calling back to the order system for every shipment.
- SaaS billing: a “usage recorded” event carrying enough tenant and plan detail for a billing service to calculate charges without querying a tenant-management service per event, at high event volume.
- Logistics: a “shipment delayed” event notification (thin) is usually sufficient, since delay handling is infrequent enough that interested services calling back for full shipment detail is cheap.
Prerequisites:
- A clear, deliberate answer — per event type, not once for the whole system — to which shape fits, based on consumer volume, data size, and staleness tolerance, rather than a blanket policy.
- Consumers building local read models from fat events need their own reconciliation strategy for missed or out-of-order messages (Section 7).
4. When Not to Use It
- Using event-carried state transfer as a default for every event, “to avoid callbacks.” This is the actual failure Northwind almost had with
analytics-pipeline’s full-payload request — every fieldorder-serviceever needs to change becomes a negotiation with every consumer that embedded it, defeating the loose coupling Chapter 7 established. - Using thin notification for a high-volume, latency-sensitive consumer. Forcing
loyalty-serviceto call back toorder-servicefor every single order recreates a synchronous dependency at exactly the volume where Chapter 7’s cascading-failure risk reappears. - Risk of both, done carelessly: a fat event whose embedded data silently becomes the de facto contract even for fields the publisher considered incidental — a consumer built against
OrderPlaced.customerTierfor a one-off report now depends on a field the publisher didn’t intend to promise long-term. Document embedded fields as explicitly as any other API contract, whichever shape you choose.
5. Implementation Example
analytics-pipeline — a low-frequency, exploratory consumer — uses thin notification and calls back:
package `in`.o612.eng.northwind.order.api
import java.util.UUID
/** Thin notification — the only guaranteed contract is that this order * exists and this fact happened. Consumers needing more call GET /orders/{id}. */data class OrderPlacedNotification(val orderId: UUID, val occurredAt: java.time.Instant)package `in`.o612.eng.northwind.analytics
import `in`.o612.eng.northwind.order.api.OrderPlacedNotificationimport org.springframework.kafka.annotation.KafkaListenerimport org.springframework.stereotype.Component
@Componentclass OrderPlacedListener(private val orderServiceClient: OrderServiceClient, private val warehouse: AnalyticsWarehouse) {
@KafkaListener(topics = ["orders.notifications"], groupId = "analytics-pipeline") fun onOrderPlaced(notification: OrderPlacedNotification) { // Low volume, exploratory use — a callback per event is an // acceptable, explicit trade-off here, not an accident. val fullOrder = orderServiceClient.getOrder(notification.orderId) warehouse.recordOrder(fullOrder) }}loyalty-service — high-volume, needs tier data on every event — uses event-carried state transfer and never calls back:
data class OrderPlacedForLoyalty( val orderId: UUID, val customerId: UUID, val customerTier: String, // GOLD | SILVER | STANDARD, snapshotted at order time val totalAmount: java.math.BigDecimal, val occurredAt: java.time.Instant,)package `in`.o612.eng.northwind.order.internal
import `in`.o612.eng.northwind.order.api.OrderPlacedForLoyaltyimport `in`.o612.eng.northwind.order.api.OrderPlacedNotificationimport org.springframework.kafka.core.KafkaTemplateimport org.springframework.stereotype.Component
@Componentinternal class OrderEventPublisher( private val kafkaTemplate: KafkaTemplate<String, Any>,) { fun publishNotification(order: Order) { kafkaTemplate.send("orders.notifications", order.id.toString(), OrderPlacedNotification(order.id, order.placedAt)) }
fun publishForLoyalty(order: Order, customerTier: String) { kafkaTemplate.send("orders.loyalty-events", order.id.toString(), OrderPlacedForLoyalty(order.id, order.customerId, customerTier, order.totalAmount, order.placedAt)) }}Two separate topics, deliberately — orders.notifications (thin, broad audience) and orders.loyalty-events (fat, narrow audience) — rather than one topic everyone subscribes to and filters. This keeps each topic’s schema evolution independent: analytics-pipeline’s topic can add fields without loyalty-service’s consumer ever noticing, and vice versa.
package `in`.o612.eng.northwind.loyalty
import `in`.o612.eng.northwind.order.api.OrderPlacedForLoyaltyimport org.springframework.kafka.annotation.KafkaListenerimport org.springframework.stereotype.Component
@Componentclass OrderPlacedForLoyaltyListener(private val pointsCalculator: PointsCalculator) {
@KafkaListener(topics = ["orders.loyalty-events"], groupId = "loyalty-service") fun onOrderPlaced(event: OrderPlacedForLoyalty) { // No callback to order-service at all — every field needed is // already in the event, snapshotted at the moment of ordering. pointsCalculator.awardPoints(event.customerId, event.customerTier, event.totalAmount) }}Reconciliation job, guarding against the staleness/drift risk that comes with any consumer maintaining its own copy of data from events:
package `in`.o612.eng.northwind.loyalty
import org.springframework.scheduling.annotation.Scheduledimport org.springframework.stereotype.Component
@Componentclass TierDriftReconciliationJob( private val customerServiceClient: CustomerServiceClient, private val loyaltyRepository: LoyaltyRepository,) { @Scheduled(cron = "0 0 3 * * *") fun reconcileTierSnapshots() { // Nightly, low-frequency check: compare a sample of recently-used // tier snapshots against the current source of truth, and alert // if drift exceeds a threshold — a safety net, not the primary path. val driftCount = loyaltyRepository.recentTierSnapshots() .count { snapshot -> customerServiceClient.currentTier(snapshot.customerId) != snapshot.tier } if (driftCount > DRIFT_ALERT_THRESHOLD) { // emit an alert metric — see Section 7 } }
companion object { const val DRIFT_ALERT_THRESHOLD = 50 }}6. Step-by-Step Flow
- Client action. Unchanged from every previous chapter —
POST /orders. - API request/event publication.
order-servicepublishes two distinct events to two distinct topics after commit, each shaped for its audience. - Service behavior.
analytics-pipelinereacts to the thin notification by calling back;loyalty-servicereacts to the fat event with everything it needs already present. - Database interaction. Each consumer writes to its own database, per the database-per-service discipline from Chapter 3 —
loyalty-service’s points ledger andanalytics-pipeline’s warehouse are both independently owned. - Inter-service communication.
analytics-pipeline’s callback is the only synchronous hop in this entire flow — a deliberate, documented exception, not an accident. - Error or failure handling. If
order-serviceis down whenanalytics-pipeline’s callback fires, that specific event’s processing fails and retries later (dead-letter handling is the next chapters’ subject) — a costloyalty-servicenever pays, since it needs nothing further fromorder-serviceto finish processing. - Observability signals. Track callback rate and callback failure rate for
analytics-pipelinespecifically — a rising callback failure rate is a direct, measurable cost of the thin-notification choice for that consumer. - Final response/outcome. Both consumers end up with correct data, using two different mechanisms chosen deliberately for their two different volume and coupling profiles.
7. Production Concerns
- Timeouts, retries, idempotency. The callback path (thin notification) needs its own retry and timeout handling, exactly like any synchronous call (Chapter 6) — it hasn’t stopped being a network call just because it started life as a reaction to an event.
- Data consistency and transaction boundaries. Fat events snapshot data at publish time —
customerTierinOrderPlacedForLoyaltyreflects the tier when the order was placed, not the tier right now. This must be a documented, intentional semantic, not an accidental staleness bug — for loyalty-point calculation, “tier at time of purchase” is usually the correct business rule, not a compromise. - Schema evolution and backward compatibility. Fat events accumulate more fields that more consumers depend on over time — track which consumers use which fields (a schema registry with compatibility checking helps concretely here) so a field removal doesn’t silently break a consumer nobody remembered was using it.
- API versioning. The callback endpoint (
GET /orders/{id}) used by thin-notification consumers is now a real, ongoing API contract, not an implementation detail — version and document it with the same care as any other public REST contract from Chapter 2 onward. - Authentication and service-to-service trust. Callback calls need the same service-to-service authentication as any other synchronous call (Chapter 6’s concerns apply in full) — an event’s origin doesn’t grant the resulting callback any special trust.
- Logging, metrics, tracing, correlation IDs. Propagate the original correlation ID through both the event and any resulting callback, so a trace of “what happened to order X” includes the callback hop, not just the event publish.
- Kubernetes deployment. No special considerations beyond Chapter 7’s — but note that a thin-notification consumer’s callback adds load back onto the publisher (
order-service), which must be capacity-planned for that traffic, not just for its own direct client traffic. - Testing strategy. Test fat-event consumers by asserting behavior purely from the event payload (no network mocking needed, since there’s no callback) — a genuinely simpler test than a thin-notification consumer, which needs its callback mocked (WireMock, as in Chapter 2) in addition to the event itself.
- Migration strategy. If a thin-notification consumer’s callback volume grows to the point where it’s a real load or availability concern, migrate that specific consumer to a fat event (or a dedicated topic) — a decision made per consumer, based on its actual growth, not a system-wide redesign.
8. Common Mistakes
- One giant event trying to serve every consumer. Cramming everything every current and hypothetical future consumer might need into a single
OrderPlacedschema makes every field a negotiation across every subscriber. Fix: split into purpose-specific topics/events per consumer need, as this chapter does withorders.notificationsandorders.loyalty-events. - Defaulting to thin notification for a high-volume consumer. Forcing a per-event callback for a consumer processing thousands of events per second reintroduces the exact temporal coupling event-driven architecture was adopted to remove. Fix: measure consumer call volume and choose the shape per consumer, not by a blanket policy.
- Treating an embedded field’s staleness as a bug instead of a documented semantic. Alerting on
customerTierinOrderPlacedForLoyalty“not matching the customer’s current tier” without recognizing that “tier at time of order” is the intended business meaning wastes on-call time chasing a non-issue. Fix: document the intended point-in-time semantics of every embedded field explicitly, in the event’s own schema comments. - No reconciliation for consumers building local read models from fat events. Assuming every event will always be delivered, in order, forever, without a periodic sanity check invites silent, undetected drift. Fix: build a low-frequency reconciliation job (as shown in Section 5) comparing a sample of the local copy against the source of truth.
- Letting a callback-based consumer’s failure block the event pipeline. If
analytics-pipeline’s consumer thread blocks retrying a failed callback without a bound, it can stall processing of subsequent messages on its partition. Fix: bound callback retries, and route persistently failing messages to a dead-letter topic (the next chapters address this explicitly) rather than blocking the consumer indefinitely. - Not treating a widely-embedded field as a public contract. Removing or renaming a field from a fat event because “it seemed unused” without checking which consumers actually depend on it breaks them silently, discovered only when their local read model stops updating correctly. Fix: treat every field in a widely-consumed event exactly as seriously as a REST API field (Chapter 2’s contract discipline applies here too).
9. Decision Guide
| Problem signal | Use this pattern? | Why | Alternative |
|---|---|---|---|
| Consumer needs data on every event, at high volume | Event-carried state transfer | Avoids reintroducing a synchronous dependency at exactly the volume that matters most | — |
| Consumer needs extra data rarely, or the data is large/volatile | Event notification (thin) + callback | Keeps the common-path event small; pays the callback cost only when actually needed | — |
| Many consumers, each needing different subsets of data | Multiple topics, shaped per consumer | Keeps each schema’s evolution independent and each payload minimal for its audience | One large shared event (avoid) |
| Embedded data’s staleness is unacceptable for the consumer’s use case | Event notification (thin) + callback | A callback gets current data; an embedded snapshot may be stale by design | — |
| Data volume/frequency is low and coupling isn’t a proven issue yet | Either — start thin | Simpler to implement and evolve; upgrade to a fat event only if callback load becomes a real cost | — |
10. Hands-On Exercise
Extend it: design the event shape for a new fraud-check consumer that needs to score every order in real time against the customer’s recent order history and payment method — does this consumer’s profile point toward thin notification or event-carried state transfer? Justify using Section 3’s criteria, then implement the chosen shape.
Simulate a failure: have order-service stop responding to GET /orders/{id} entirely (simulating an outage) while analytics-pipeline’s callback-based consumer is processing a backlog of notifications. Observe what happens to that consumer’s processing, and design a dead-letter or retry-with-backoff strategy so a order-service outage doesn’t stall analytics-pipeline indefinitely.
Decision question, with justification required: loyalty-service’s OrderPlacedForLoyalty event currently embeds customerTier as a plain string. A new requirement: loyalty points should also depend on the customer’s lifetime order count, a value that changes on every order and would need to be embedded fresh each time, doubling as a potential source of race conditions if two orders are placed nearly simultaneously. Should this value be embedded in the event, fetched via callback, or computed independently inside loyalty-service from its own local event history? Weigh the trade-offs from Section 1 and Section 7.
11. Key Takeaways
- Publish-subscribe (Chapter 7) answers who gets a copy of an event; event notification versus event-carried state transfer answers how much each event carries — a separate, per-event-type decision.
- Thin notifications minimize schema coupling but may reintroduce a synchronous callback dependency — acceptable for low-volume, exploratory, or staleness-sensitive consumers.
- Fat events (event-carried state transfer) eliminate the callback but widen the schema’s blast radius — justified for high-volume consumers that would otherwise recreate the temporal coupling event-driven architecture was meant to remove.
- Different consumers of the same underlying fact can and often should receive differently-shaped events on different topics, rather than one universal event trying to serve everyone.
- Embedded data in a fat event is a point-in-time snapshot, not a live value — document that semantic explicitly so it isn’t mistaken for a staleness bug.
- Any consumer building a local read model from fat events needs a periodic reconciliation check against the source of truth — events can be missed or misordered, and silent drift is worse than a visible failure.
- Treat every field in a widely-consumed event with the same contract discipline as a public REST API field — removing one breaks consumers you may not know exist.