Fan one AtlasMart change stream into independent derived stores, expose freshness/version evidence, inject projection failure, and prove rebuild/reconciliation behavior.
Build Search Indexes, Caches, Warehouses, and Read Models from Change Streams
One source stream can feed many read-optimized stores, but each projection becomes an independently failing operational copy.
Build multiple independent derived stores from one retained change stream and record source version/offset in every projection.
Expose projection freshness and diagnose partial consumer failure rather than treating all derived stores as synchronously current.
Design replay/backfill/rebuild procedures and schema-evolution compatibility for search, cache, warehouse, and CQRS-style read models.
Give each projection independent observability, reconciliation, security, and recovery ownership.
1. Derived stores are operational copies, not free indexes
AtlasMart wants a search index for catalog text, a low-latency product cache, an analytics warehouse, and a read model tailored to the storefront. One authoritative catalog/order stream can feed all four. Each target is a projection: a representation optimized for a different read workload. The same source event can therefore fan out into four writes with four different failure modes and freshness profiles.
The design becomes manageable when every projection records the source identity it has applied—an offset, source version, event ID, or combination. Without provenance, “is this stale?” becomes guesswork.
2. Independent checkpoints and freshness evidence
| Projection | Optimization | Useful freshness evidence | Typical rebuild source |
|---|---|---|---|
| Search | text/vector retrieval | source version + stream offset + indexing lag | retained events or authoritative catalog snapshot |
| Cache | low-latency key lookup | source version + age/TTL | authoritative API/database |
| Warehouse | large scans/analytics | ingestion watermark + source partition offsets | change history plus periodic snapshot |
| CQRS read model | screen/API query shape | aggregate sequence/version | event log/change stream |
Freshness can be part of the API. For example, an inventory hint
may show source_version=22 while the authoritative
inventory version is 24. A business-critical purchase decision
should not infer stock correctness from that hint.
3. Partial consumer failure must remain local and diagnosable
Suppose the search consumer misses catalog offset 2 while cache, warehouse, and read-model consumers continue. Search still returns the old title, but the others are current. Calling “the pipeline healthy” because three of four consumers are green hides the exact failure the architecture introduced.
Each target therefore needs its own last-success time, source offset/version, error/retry counts, DLQ/quarantine, schema version, and reconciliation status. Failure isolation is a benefit only if operations can see the isolated failure.
4. Deliberately wrong approach: make projections authoritative because they are convenient
If customer service edits inventory directly in the search index because that UI is easier to access, the projection now has an independent write authority. Rebuild from the source destroys the manual change; failing to rebuild leaves undocumented divergence. Derived stores should usually be read-only from the business perspective, with repair flowing from the authoritative owner.
Schema evolution also belongs in the projection contract. Additive fields can use defaults for older events; incompatible meaning changes require a new schema/version and possibly a backfill. Do not let a connector-generated payload silently become an eternal public API.
5. AtlasMart lab: fan-out, stale search, rebuild, schema compatibility
Python 3.13+ standard library only. All stores are dictionaries/lists. No search engine, cache server, warehouse, CDC connector, broker, or cloud service is required.
events = [
{"offset":0,"sku":"sku-7","version":1,"title":"Atlas Lamp","price":100,"stock":5},
{"offset":1,"sku":"sku-7","version":2,"title":"Atlas Lamp","price":120,"stock":5},
{"offset":2,"sku":"sku-7","version":3,"title":"Atlas Lamp Pro","price":120,"stock":4},
]
search={}; cache={}; warehouse=[]; read_model={}
def apply_search(e): search[e['sku']]={"title":e['title'],"source_version":e['version'],"offset":e['offset']}
def apply_cache(e): cache[e['sku']]={"price":e['price'],"stock":e['stock'],"source_version":e['version'],"offset":e['offset']}
def apply_wh(e): warehouse.append({"sku":e['sku'],"version":e['version'],"price":e['price'],"offset":e['offset']})
def apply_read(e): read_model[e['sku']]={"title":e['title'],"price":e['price'],"stock":e['stock'],"source_version":e['version']}
print("FAN-OUT WITH ONE CONSUMER FAILURE")
for e in events:
if e['offset'] != 2: # search consumer misses offset 2
apply_search(e)
apply_cache(e); apply_wh(e); apply_read(e)
print("search:", search['sku-7'])
print("cache:", cache['sku-7'])
print("read model:", read_model['sku-7'])
print("warehouse rows:", len(warehouse))
source_version=events[-1]['version']
for name,store_version in [
("search",search['sku-7']['source_version']),
("cache",cache['sku-7']['source_version']),
("read_model",read_model['sku-7']['source_version'])]:
print(name,"freshness lag in versions:",source_version-store_version)
print("\nREBUILD SEARCH FROM RETAINED CHANGE HISTORY")
search.clear()
for e in events: apply_search(e)
print("rebuilt search:", search['sku-7'])
print("\nSCHEMA EVOLUTION")
new_event={"offset":3,"sku":"sku-8","version":1,"title":"Desk","price":50,"stock":7,"currency":"USD"}
def compatible_price(e):
return {"amount":e['price'],"currency":e.get('currency','USD')}
print("backward-compatible reader:", compatible_price(events[0]), compatible_price(new_event))
print("each projection owns independent checkpoints, freshness SLOs, rebuild logic, and reconciliation")
The search projection remains at source version 2 while
cache/read model reach version 3. The lab then clears search
entirely and rebuilds it from retained history, proving
disposability. An added currency field is handled
by a backward-compatible reader.
6. Production judgment
Define a freshness service-level objective (SLO) per projection. A cache might tolerate seconds; fraud or authorization projections may require much tighter boundaries; a warehouse may accept minutes. Separate “event successfully consumed” from “target durably indexed and queryable.” Track ingest-to-visible latency and not just consumer offset. Protect sensitive fields during fan-out: the warehouse, search index, and cache should receive only the data necessary for their purpose.
Rebuild plans need source retention or a snapshot plus replay boundary, capacity headroom, throttling, validation, cutover, and rollback. A full rebuild often generates far more I/O than steady-state change application.
The next lesson adds continuous reconciliation because replayable streams do not prove derived state stayed correct forever.
Check your understanding
- Why should a projection store source version or offset?
- Can one healthy consumer imply all derived stores are current?
- Why should a derived store usually not accept authoritative writes?
- What is required for a safe full rebuild?
- Why is consumer offset not the same as query freshness?
Review the answers
1. It makes freshness, replay progress, drift detection, and repair decisions observable.
2. No. Each projection has independent checkpoints, target failures, indexing visibility, and schema handling.
3. Those writes are not represented in the source and can be lost on rebuild or create permanent multi-writer divergence.
4. A trusted source/snapshot, replay boundary, capacity/throttling plan, validation, cutover, and rollback.
5. The target may acknowledge ingestion before data is durably indexed or visible to queries.
References
Foundational claims use primary specifications/research or current official documentation where practical. Product references are optional implementation anchors; the mandatory labs are vendor-neutral.
- PostgreSQL 18 logical decoding — Authoritative source-stream example for building external consumers.
- Debezium documentation — Official CDC connector documentation and event formats for optional implementation study.
- Apache Kafka design documentation — Retained-log and replay concepts for optional stream implementation.
- Debezium 3.6.1.Final release notes — Current stable version snapshot and data-integrity/security fixes.