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.

Advanced120–155 minutesDerived-projection labPython 3.13+ · standard library / sqlite3 where notedVendor-neutral · free/local mandatory pathLast reviewed: August 2026
01

Build multiple independent derived stores from one retained change stream and record source version/offset in every projection.

02

Expose projection freshness and diagnose partial consumer failure rather than treating all derived stores as synchronously current.

03

Design replay/backfill/rebuild procedures and schema-evolution compatibility for search, cache, warehouse, and CQRS-style read models.

04

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

Mandatory lab environment

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.

python · AtlasMart deterministic simulation
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")
Expected evidence

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

  1. Why should a projection store source version or offset?
  2. Can one healthy consumer imply all derived stores are current?
  3. Why should a derived store usually not accept authoritative writes?
  4. What is required for a safe full rebuild?
  5. 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.

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.