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.

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.