ddia/ch12.md - DDIA
Reason about data flows
- Writing to 2 locations introduces inconsistencies. Better derive from event logs. Event logs can be made idempotent.
- Fundamentally, a streaming framework from primary to derived almost always should be:
- fault-tolerant -- do not lose any data. this causes derived states to become out of sync
- ordered -- to keep data consistent and often just correct, events must be arriving in order.
Problem with logical clocks
- Scalar logical clock: gives total order, but no info about concurrency or conflict.
- Vector clocks:
- message size scales linearly with amount of replicas
- doesn't guarantee serializability or linearizability which is often wanted to systems where the consistency matters, making vector clocks "middle" ground unused for both consistent, and available systems.
- you will still need to deal with conflicting writes (concurrent writes become conflicting).
Batch vs. stream processing
- When to use which?
- The main fundamental difference is that stream processors operate on unbounded datasets, whereas batch process inputs are of a known, finite size.
- Batch problems
- latency.
- also, most systems assume you run over big chunks
- can't be ran too often -- if you do run too often, you pay the fixed overhead per-run
- after all, it seems you are just fighting the model.
- latency.
Replicas, indices, etc.
- Setting up a database index, new replica in distributed system, or replaying the event logs -- are ALL the special cases of setting up a new derived system.
Pushing state to clients
- How to deal with client outages in push-based model?
- Maintain a client offset on data that the user can later query.
Distributed query execution
- We can reuse stream processing to perform distributed query execution
- Use-case: perform read over sharded data
- Get the event to the source of stream, then fan-out to different shards, collect data, push down to the next stage where this data is joined together and returned back
Distributed global constraints
- For example, enforcing a uniqueness constraint requires consensus.
- If we have a global constraint (invariant) that the system must maintain, we have to coordinate at all times.
- If it's acceptable for the system, sometimes we can cheat by not coordinating and only preventing the user from seeing the broken invariant (non-negative global counter violated => just clamp it), or performing a compensating operation (as far as sending sorry emails that your operation is actually cancelled, although the system ACKed it previously).
Consistency
- By "consistency" we often conflate 2 things:
- timeliness
- integrity
- The CAP theorem uses “consistency” in the sense of linearizability, which is a strong way of achieving timeliness
Coordination avoidance
- No coordination almost always means better performance, therefore it makes sense to reduce it. Oftentimes, we can assume loose constraints and performing compensating operations to avoid coordination at write-time.