MillWheel: Fault-Tolerant Stream Processing at Internet Scale - Google, 2013
Paper: PDF
Mirror: PDF
Data Model
- Each record is a
(key, value, timestamp)tuple - The key is used to partition across instances of a computation.
- Each record is processed separately.
- For optimization, side effects of computation over several records can be batched.
Streams
- Streams are streams of records. Different computation nodes care about different data from streams.
Computations
- Computations are application logic that lives in nodes. Computation code is invoked upon receipt of input data, at which point user-defined actions are triggered
- Computation nodes process records based on a specific key. This means that the volume of a stream is partitioned by key and can be therefore scaled.
Persistent state
- Problem: computation node can crash and loose all data it accumulated/read from the stream.
- Solution: have some persistent, fault-tolerant state.
- "Persistent state is an opaque byte string that is managed on a per-key basis"
- The user provides serialization and deserialization routines (such as translating a rich data structure in and out of its wire format), like protobuf
- Persistent state is backed by a replicated, highly available data store (like Bigtable or Spanner), which ensures data integrity in a way that is completely transparent to the end user.
- The check pointing is atomic per key. This allows exactly-ones semantics.
Watermarking
- Each record coming in to the system is given a timestamp (often corresponding to the time the data was generated)
- The system tells the computations that all the data up to the low watermark has likely already arrived. It just recursively takes the min of the low watermark at input streams. The base case is the lowest unprocessed timestamp.
low watermark of A = min(oldest work of A, low watermark of C : C outputs to A)
- Per computation node, I want a single number "what is the oldest timestamp I still can" => Everything earlier can be processed.
- This leads to a problem that 1 straggler computation node can lead to hang of the whole pipeline. Possible solution: advance watermark regardless (99% of records arrived).
Exactly-ones semantics
- Exactly-ones semantics is either:
- at-least-ones semantics + deduplication.
- at-least-ones semantics + idempotency.
- Process for exactly once delivery:
- Record is checked against dedupe data
- User code is run
- Pending changes are committed
- Senders ACKed
- Productions sent to downstream services
Zombies
- Because computation nodes can be replicated, there is a danger of zombie writers
- Classic solution: attach a sequence token to each write.
Load balancing
- Master divides the key space in intervals and distributes the responsibility for intervals across the computation nodes.
- Intervals can be moved/split/merged as a response to changing load.
- Each interval has a unique sequencer, invalidated on interval change
Types of productions (messages to downstream)
- Production = message to a downstream computation
- There 2 kinds of productions: strong and weak
- Strong productions are used to guarantee exactly-once semantics (required to non-idempotent operations)
- They are based on sending ACK to sender only after committing to data store.
- This guarantees that if we fail before committing to data store, the sender retries and new operation is performed over the store without side-effects.
- The steps are described in "Exactly-ones semantics" above.
- Problem: latency incurred by round-trip to data store + ACK to sender.
- BUT if the underlying computation is already idempotent, we can avoid these calls by directly passing production (called weak) downstream before committing to memory.
- Then, we wait for the downstream ACK => this approach couples 2 computation nodes.