spanner

Google spanner is a scalable SQL/relational database, which is interesting because creating transactions across multiple machines is generally a hard thing to do.

  • Atomicity is guaranteed using two phase commit
  • Spanner scale by partitioning data into shards
  • External consistency means the database’s transaction order must agree with the order users can observe in the real world
    • Done by delaying messages until it is definitely in the past. The delay is reduced by using specialized hardware (GPS receivers, atomic clocks, etc)

Writing to spanner

Spanner guarantee external consistency. If a client waited for transaction 1 (T1) to complete, then initiated T2, the client must always see T1 before T2 (either T1 and T2 didn’t happen; only T1 happened; or T1 and T2 happened). How does this work?

  • When the coordinator node commits T1, it has some belief on what the current time range is [earliest1, latest1]. Importantly, the current time is between [earliest, latest]. Before telling the participants that the commit has happened (via two phase commit), the coordinator node waits until the current time range is [earliest2, latest2] definitely after the commit time. This is done by checking that earliest2 > latest1. Only after this point will T2 be able to initiate (because client waits for T1 to complete).
  • Lets use a real example:
    • client ask spanner to commit T1.
    • coordinator A runs 2PC, and get yes for everything.
    • current time of coordinator A is [100, 110]. Coordinator A commits the message with timestamp 110.
    • coordinator A waits until its current time is after 110. eg [111, 120]. Only then does it tell all participants that the commit has happened. The commit still has the same timestamp of 110.
    • client confirms that T1 has committed, and now ask spanner to commit T2.
    • coordinator B runs 2PC, and get yes for everything.
    • current time of coordinator B is [120, 130]. Coordinator B commits the message with timestamp 130.
      • no matter what, latest time of coordinator B cannot be 110, because it is guaranteed that “true time” is after time 110 (coordinator A has waited for time 110 to past).
    • coordinator B waits until its current time is after 130. eg [140, 150]. Only then does it tell all participants that the commit has happened. The commit still has the same timestamp of 130.
  • if an external host query spanner at the following transaction time:
  • before 110: no host will say they have seen T1 because T1 has a timestamp of 110. no host will say they have seen T2 because T2 has a timestamp of 130.
  • [110, 130): all host will say they have seen T1 because T1 has a timestamp of 110. no host will say they have seen T2 because T2 has a timestamp of 130.
  • after and include 130: all host will say they have seen T1 because T1 has a timestamp of 110. all host will say they have seen T2 because T2 has a timestamp of 130.

Reading from spanner

When a client reads from spanner, spanner chooses a read timestamp. A few things must be true before spanner responds to the request:

  • read timestamp is chosen such that any other host using this reading as a “happened-before” event can observe this effect. the obvious choice is to use TrueTime.latest (upperbound of truetime).
  • all events/transactions that happened before this read timestamp must be seen by the host servicing this read request. This ensures that asking two different host for the same data given the same read timetamp don’t yield inconsistent result. This is done by the host confirming with the paxos leader that it has receives all messages up to read timestamp.

Why is truetime useful?

  • TrueTime returns a range [earliest, latest] which is guaranteed to contain the actual real-world time.
  • Because this guarantee is available on every Spanner machine, independent coordinators can assign globally meaningful timestamps without contacting a centralized timestamp server.
  • A write chooses a commit timestamp >= TT.now().latest and does not expose the transaction as committed until TT.now().earliest is greater than that commit timestamp.
  • The GPS receivers and atomic clocks used by TrueTime keep the uncertainty interval small. A smaller uncertainty interval means a shorter commit-wait