Real-Time Analytics Platform  ·  View 11 of 21  ·  Runtime

Stream Processing Pipeline

What each micro-batch does, stage by stage, and exactly where the delivery guarantee is established.

Editable source SVG draw.io All views
Read
Read
Checkpoint & Offsets
advanced after sink ack
Checkpoint & Offsets...
Event Hubs Source
maxOffsetsPerTrigger
Event Hubs Source...
Backpressure Control
trigger 1 s
Backpressure Control...
Decode
Decode
Undecodable Event
Undecodable Event
Avro Decode
Avro Decode
Schema Resolve
by schema ID
Schema Resolve...
Enrich
Enrich
Stateful Lookup
RocksDB state store
Stateful Lookup...
Broadcast Dimension Join
Delta reference
Broadcast Dimension Join...
Dimension Miss
late reference fallback
Dimension Miss...
Window
Window
Watermark 2 min
Watermark 2 min
Window Assignment
tumbling · sliding · session
Window Assignment...
Stateful Aggregation
per key and window
Stateful Aggregation...
Write
Write
Delta Silver Sink
txnAppId idempotence
Delta Silver Sink...
ADX Streaming Sink
ingest-by dedup tag
ADX Streaming Sink...
Dead-Letter Hub
quarantine · replayable
Dead-Letter Hub...
Signal
Signal
Gold Event Hub
enriched output
Gold Event Hub...
Checkpoint Commit
after sink ack
Checkpoint Commit...
Pipeline Metrics
lag · duration · rows
Pipeline Metrics...
decode failure
decode failure
beyond lateness
beyond lateness
Stream Processing Pipeline
Stream Processing Pipeline
Data store
Data store
Queue / topic
Queue / topic
Security / platform
Security / platform
Risk / gap
Risk / gap
Application we own
Application we own
Decision point
Decision point
failure / alternate
failure / alternate
Checkpoints commit only after every sink acknowledges, giving effectively-once delivery end to end.
Checkpoints commit only after every sink acknowledges, giving effectively-once delivery end to end.
v 1.0 · owner Streaming Engineering · date 2026-08
v 1.0 · owner Streaming Engineering · date 2026-08
Text is not SVG - cannot display

Exactly-once, precisely

  • Transport is at-least-once; sinks are idempotent, giving effectively-once end to end
  • ADX deduplicates on an ingest-by tag; Delta writes carry txnAppId and txnVersion
  • The checkpoint advances only after every sink acknowledges, so a crash replays, never skips

Backpressure

  • maxOffsetsPerTrigger caps each batch so a burst lengthens the queue instead of the batch
  • Event Hubs retention is the buffer: 7 days of absorption before anything is at risk
  • Autoscale triggers on consumer lag, not on CPU, because lag is the signal that matters

Failure handling

  • Undecodable and unprocessable events are quarantined with the failure reason attached
  • A dead-letter event is never silently dropped and is replayable after the fix
  • State store corruption is recovered by replay from bronze, not by partial checkpoint repair