Order event-stream topology
Topics, consumer groups and the dead-letter path
Producers
Transit
Processors
State and recovery
Consumers
Checkout API
Order producer
Billing API
Payment producer
orders.v1
Topic
12 partitions
payments.v2
Topic
8 partitions
Order validate
Group fulfillment
Payment enrich
Group analytics
Order state
Materialized view
events.dlq
Poison events
Fulfillment
Shipping workflow
Analytics
Streaming facts
Replay tool
Approved batch
On-call
DLQ owner
OrderPlaced
PaymentCaptured
ordered
payment facts
valid order
enriched
ready orders
order facts
invalid
poison
sample
approved replay
diagram.html#The link keeps the view, selection, route and playback.
About this diagram
Data moving through stages: pipelines, ETL and ELT, change data capture, event streams, lineage.
Ask for one like it
Show how data flows through our feature platform, from the product databases to the models that read it.
The JSON
{ "kind": "dataflow", "density": "compact", "title": "Order event-stream topology", "subtitle": "Topics, consumer groups and the dead-letter path", "direction": "RIGHT", "phases": [ { "id": "producers", "label": "Producers", "nodes": ["checkout", "billing"] }, { "id": "transit", "label": "Transit", "nodes": ["orders", "payments"] }, { "id": "processors", "label": "Processors", "nodes": ["validate", "enrich"] }, { "id": "state", "label": "State and recovery", "nodes": ["order-state", "dlq"] }, { "id": "consumers", "label": "Consumers", "nodes": ["fulfillment", "analytics", "replay", "ops"] } ], "nodes": [ { "id": "checkout", "type": "service", "card": { "title": "Checkout API", "subtitle": "Order producer" } }, { "id": "billing", "type": "service", "card": { "title": "Billing API", "subtitle": "Payment producer", "brand": "stripe" } }, { "id": "orders", "type": "queue", "card": { "title": "orders.v1", "subtitle": "Topic", "tag": "12 partitions" } }, { "id": "payments", "type": "queue", "card": { "title": "payments.v2", "subtitle": "Topic", "tag": "8 partitions" } }, { "id": "validate", "type": "service", "card": { "title": "Order validate", "subtitle": "Group fulfillment" } }, { "id": "enrich", "type": "service", "card": { "title": "Payment enrich", "subtitle": "Group analytics" } }, { "id": "order-state", "type": "database", "card": { "title": "Order state", "subtitle": "Materialized view" } }, { "id": "dlq", "type": "queue", "card": { "title": "events.dlq", "subtitle": "Poison events" } }, { "id": "fulfillment", "type": "service", "card": { "title": "Fulfillment", "subtitle": "Shipping workflow" } }, { "id": "analytics", "type": "database", "card": { "title": "Analytics", "subtitle": "Streaming facts" } }, { "id": "replay", "type": "security", "card": { "title": "Replay tool", "subtitle": "Approved batch" } }, { "id": "ops", "type": "external", "card": { "title": "On-call", "subtitle": "DLQ owner", "brand": "pagerduty" } } ], "edges": [ { "id": "e1", "from": "checkout", "to": "orders", "label": "OrderPlaced", "tone": "main" }, { "id": "e2", "from": "billing", "to": "payments", "label": "PaymentCaptured", "tone": "main" }, { "id": "e3", "from": "orders", "to": "validate", "label": "ordered", "tone": "main" }, { "id": "e4", "from": "payments", "to": "enrich", "label": "payment facts", "tone": "main" }, { "id": "e5", "from": "validate", "to": "order-state", "label": "valid order", "tone": "main" }, { "id": "e6", "from": "enrich", "to": "order-state", "label": "enriched" }, { "id": "e7", "from": "order-state", "to": "fulfillment", "label": "ready orders", "tone": "main" }, { "id": "e8", "from": "order-state", "to": "analytics", "label": "order facts" }, { "id": "e9", "from": "validate", "to": "dlq", "label": "invalid", "tone": "error" }, { "id": "e10", "from": "enrich", "to": "dlq", "label": "poison", "tone": "error" }, { "id": "e11", "from": "dlq", "to": "ops", "label": "sample", "tone": "error" }, { "id": "e12", "from": "dlq", "to": "replay", "label": "approved replay", "kind": "async" } ], "notes": [ { "title": "Failure ownership", "items": [ "Poison events land in a retained dead-letter topic", "On-call inspects samples before replay", "Replay is gated, batched and auditable" ] } ]}More dataflow examples
All examples- stackmap: JSON to HTMLDataflow9 nodesWritten by an agent“Make an architecture diagram of this repository: what the packages are, what each one does, and how data moves between them from an agent-written diagram JSON to the HTML file a user opens. Back it with evidence from the code.”
- Product analyticsDataflow10 nodes
- ML feature platformDataflow13 nodesWritten by an agent“Show how data flows through our ML feature platform. Product databases (Postgres) are captured with Debezium CDC into Kafka. A Flink job computes streaming features and writes them to Redis (online store) and to a Delta Lake on S3 (offline store). Airflow runs nightly Spark jobs over the lake to build training sets, which a training pipeline on SageMaker consumes; models are registered in MLflow. The prediction service reads online features from Redis and loads models from MLflow. Also note that raw Kafka topics are archived to S3 for replay.”