flowchart LR
src1[("Orders DB")] --> cdc["CDC connector"]
cdc --> t1[["orders.raw<br/><i>12 parts · key: orderId<br/>retain 7d</i>"]]
src2["Clickstream"] --> t2[["events.clicks<br/><i>36 parts · key: sessionId<br/>retain 3d</i>"]]
t1 --> p1["Normalise<br/><i>stateless</i>"]
p1 --> t3[["orders.clean<br/><i>12 parts · key: orderId<br/>compacted</i>"]]
t3 --> p2["Enrich + join<br/><i>stateful · 30 min window</i>"]
t2 --> p2
p2 --> st[("RocksDB state<br/><i>checkpoint 60s</i>")]
p2 --> t4[["orders.enriched<br/><i>12 parts · retain 30d</i>"]]
p2 --> dlq[["orders.dlq"]]
t4 --> sink1["Serving store"]
t4 --> sink2["Lakehouse sink<br/><i>5 min commit</i>"]
Data View
Streaming Topology Diagram
Topics, partitions, processors and state stores in one picture, with the retention and keying decisions that determine whether it can be replayed.
Stream Processing Frameworks
Design