Photon: Fault-tolerant and Scalable Joining of Continuous Data Streams - Google, 2013
Paper: PDF
Mirror: PDF
System overview
- We have 2 streams of continuous data and we want to join events there based on certain key
- Requirements: exactly-ones processing, good fault-tolerance
- Challenges: unordered events, delayed arrival of primary stream.
- Communication between the services is done via RPC.
- with throttling.
Why not just add all metadata to click?
- This way we could avoid the need in stream joining
- To add metadata to click we have 2 options:
- add it to URL (this is a very bad decision: security, latency (incurred by de/serialization), URL length limit)
- add to the body.
- first problem: not idiomatic for GET request.
- second: not secure. we have to store full query metadata on client -- they can inspect it
- To add metadata to click we have 2 options:
Streams
- Streams are coming in a form of logs -- stored in GFS
- Streams are unordered, arbitrary delayed.
- Each event has a key
event_id = (ServerIP:ProcessID:Timestamp)which is a unique identifier thanks to the monotonic clock.
Guarantees (imposed by usage -- advertisers billing)
- At-most ones at any time
- As near to real-time as possible
- Exactly-ones eventually
Fault-tolerance
- Revenue generating system => must survive any disaster
- Solution: the main data store must achieve consensus (using Paxos) across at least 2 DCs (data center) for any write.
- This main store is called IdRegistry
IdRegistry
- Data store for storing processed
event_ids- Before writing to the output, first store the processed
event_idthere - For cases, when we store the processed
event_idbut fail before writing to the output, we have a "consistency" background process that checks if all events from IdRegistry have been written to the output.
- Before writing to the output, first store the processed
- This system is replicated by Paxos across DCs
- This means (at least) 1 network round-trip. Assume it's 100ms, then we can have only 10 transactions a second.
- Optimizations
- shard the IdRegistry by
event_id. they use consistent hashing to optimize resharding. - server-batch writes.
- shard the IdRegistry by
- To not overwhelm IdRegistry, we throttle the requests there.
Dispatcher and Joiner
- Dispatcher reads the logs from local GFS cell.
- it uses many worker processes concurrently process the log files. The state is preserved in IdRegistry.
- before sending the event to Joiner, it deduplicates the record using IdRegistry.
- If Joiner answers with error, it exponentially retries the event.
- Joiner is the component that takes the joined event from the LogEventStore (index over primary stream), joins it with the foreign event, stores the result to IdRegistry and then to output (also GFS log).
- It's stateless, so easy to scale.
- There could be an "adapter" inside that does some filtering logic over the joined event.
LogEventStore
- Sparse index
(event_id --> offset in log file)over primary logs. Populated every T seconds, or B bytes -- whatever comes first. - Has an optimization -- CacheEventStore in front, which is a LRU cache.
- Stored in BigTable KV store (= distributed hash table).
Throttling
- There are 3 points where throttling is used.
- Dispatcher sends event to Joiner
- Joiner checks if it can accept the request and potentially rejects it.
- This helps to not overwhelm the IdRegistry too, not only Joiner. One greedy Joiner can take down the shared resource, leading to outage
- Dispatcher retry
- Caused by above; exponential backoff to avoid retry storm.
- Joiner's outstanding RPCs limit to IdRegistry.
- the unlimited amount of outstanding RPCs could inflate the amount of orphaned events.
Learnings
- Paxos writes are expensive. Minimize operations done there.
- Use throttling
High-level overview:
