ddia/ch08.md - DDIA

Anomalies in databases

  • Isolation levels defined by SQL standard. Image taken from "PostgreSQL 17 Internals" by Egor Rogov:** Isolation levels

  • My own table of anomalies in databases. X-axis has anomalies, y-axis contains isolation levels. Each cell contains a yes/no answer to possibility of anomaly. Anomalies in DBs

  • Hierarchy of consistency levels from jepsen.io: Anomalies in DBs

Details on isolation levels:

  • Read Committed
    • Практический вывод: в транзакции нельзя принимать решения на основании данных, прочитанных предыдущим оператором, ведь за время между выполнением операторов все может измениться
    • How to work around the non-repeatable read:
      • Add predicates like ALTER TABLE accounts ADD CHECK amount >= 0;
      • Use 1 operator (they are atomic): INSERT, UPDATE, DELETE, INSERT ON CONFLICT
      • Block rows using SELECT FOR UPDATE, or table using LOCK TABLE
        • Bad: defeats MVCC
  • Cursor stability
    • Read Committed + guaranteed no lost update

How to implement isolation:

  • Blocking (2PL): block rows on which you operate during TX, then release
  • 2PL = concurrency protocol. guarantees serializability by following this rule: locks are applied and removed in two phases:
    1. Expanding phase: locks are acquired and no locks are released.
    2. Shrinking phase: locks are released and no locks are acquired.

Levels of isolation in Postgres:

  • Protocol of isolation based on snapshot isolation (SI)
  • No dirty reads in PG
    • Read Uncommited == Read Committed in PG
  • Repeatable read guarantees NO phantom reads too.
  • Read Committed can lose data (TODO: How?)
  • PG-specific anomalies
    • Read skew (несогласованным чтением)
      • TX B opened. Reads balance $100 of A. TX B still NOT committed.
      • TX A opened. Transferred $100 from A to B. Committed.
      • TX B continues. Reads balance of B, sees (just transferred) $100. Concludes a sum of 2 balanced $200. Not correct
      • In other words, you read row X, someone modifies row Y (which is related to X), you read row Y => they're now inconsistent.
      • How to solve? Obv, single operator!
        • But! Make sure your operator is NOT volatile (they are by default)
      • Or use Repeatable Read
    • Несогласованное чтение вместо потерянного обновления
      • Remember that nested statements ARE NOT atomic. Atomic is only single operator statement.
    • Lost update
      • Typical concurrency problem with racing writers
      • Fixed in Repeatable Read
    • Write skew
      • Parallel writes. See below
      • Possible at Repeatable Read
    • Read-only transaction anomaly
      • Leads to final system state which can not be explained by any sequence of consecutive transactions: observer’s transaction sees an impossible state, one that should never occur from the application’s perspective
      • Possible at Repeatable Read
  • В целом, идея skew заключается в том, что я оперирую на знании о состоянии таблицы, которое могло поменяться между моим чтением и записью(write skew), или чтением (read skew).
  • In Repeatable Read and Serializable, we must be able to retry transactions due to serialization errors
  • If you use Serializable for some transaction, you MUST use it for all transactions. Otherwise, behavior is Repeatable Read.
  • Repeatable Read in PG == Serializable in Oracle == Snapshot Isolation

When to use what level of isolation

  • Read Uncommitted
    • Almost never but dirty read partially helps with lost updates
  • Read Committed
  • Repeatable reads
    • Long-running queries/processes over a database (OLAP query, backup) - we need to ensure no read skews

Snapshot isolation

  • Idea: readers never block writers, writers never block readers
  • MVCC idea:
    • TX is assigned a txid
    • When TX writes, it tags the write with txid inside inserted_by and deleted_by
    • Update is internally translated into DELETE+INSERT
    • All versions of a row are stored on the same DB heap. Versions of the same row form a linked list, sorted by age.
  • Visibility rules for consistent snapshot:
    • At start of TX, collect a list of all txids running. Ignore any writes from them.
    • New txids are ignored too
    • Write by aborted txs are ignored too
  • How to indexes work in MVCC?
    • Leaf contain left/right pointers to older/newer versions. You must iterate to find what's visible to you.

How to submit transactions from code?

  • Stored procedures
    • Bad: hard to test, manage, different languages for different databases, etc.
    • Good: less network communication
    • Some databases also have stored procedures but in general-purpose PLs: Redis uses Lua, VoltDB uses Java/Groovy, MongoDB uses JS
    • Don't overuse. Don't store all of your application logic there. DB resources should be spent on queries, not app logic. By doing that you defer DB scaling which is much harder than app scaling.
  • Interactive mode
    • No, just bad. Too-too slow

Serializability implementation

  • Locks. 2PL (SS2PL)
    • 2PL != 2PC
    • 2PL provides serializable isolation
    • Writers block readers; readers block writers.
    • Algorithm
      • If transaction wants to read, it must acquire a shared lock. More transactions can have a shared lock on the same object.
      • If transaction wants to write, it must acquire a exclusive lock. Only one transaction can have an exclusive lock over an object.
      • Once transaction acquired a lock, it must hold it until the end of transaction.
    • Easily causes deadlocks
      • DB identifies => kills one of the waiting TXs
    • Predicate-based 2PL: To prevent phantom reads, we lock not over object but over predicate. Then, all new TXs must check against a list of existing predicates.
      • Bad performance
    • Index-based 2PL: like predicate-based but put a lock over an index that forms a set in which all your values are (values of your predicate are a subset of this greater set).
      • If there's no suitable index, fall back on the entire DB
  • Single-threaded execution
    • Idea: have 1 thread that executes all transactions
    • Problem: in write-heavy workloads it's bottleneck
    • Sharding. Partial solution: shard your data by CPU cores, then you can assign 1 CPU core per shard. It assumes you don't need cross-shard transactions; otherwise, gg perf.
  • Serializable snapshot isolation (SSI)
    • Optimistic concurrency control
    • Check happen at commit-time.
    • Helps even more if you operations (TXs) are commutative
    • How to detect stale reads?
      • Keep track of requested tx versions. If they are updated => abort
      • Stale reads => Check whether ignored writes (dirty writes) have been committed

How to inspect Postgres

  • To see current locks use pg_locks table
    • Tip: join with pg_stat_activity on (pid) to match query to lock
  • pg_stat_activity to see current operations
  • heap_page_items(get_raw_page([TABLE NAME: str], [PAGE INDEX: int]) to see contents of a specific page
  • pgstattuple([TABLE NAME: str] to see stats on dead tuples (helps for vacuuming stats?)
    • needs pginspect extension.

References:

  • In-depth article about PG anomalies
    • https://dansvetlov.me/postgres-anomalies/
  • The part of PG we hate the most (Article about internals of MVCC in PG)
    • https://www.cs.cmu.edu/~pavlo/blog/2023/04/the-part-of-postgresql-we-hate-the-most.html
  • Nine ways to shoot yourself in the foot with PostgreSQL
    • https://philbooth.me/blog/nine-ways-to-shoot-yourself-in-the-foot-with-postgresql
  • What is state-machine replication?

Distributed transactions

2PC:

  • Coordinator node processes a transactions in 2 phases - PREPARE and COMMIT/ABORT.
  • When participant answers "yes" to PREPARE, participant cannot abort.
  • When coordinator receives only "yes" from all participants, it decides to saves COMMIT to its log.
  • Coordinator sends COMMIT.
  • Participant waits for COMMIT or ABORT. If no requests incomes, it waits.

3PC:

  • Assumes bounded network delays and bounded response times to guarantee atomicity.

Better solution than 2PC or 3PC is fault-tolerant consensus protocol.

Message processing

  • Exactly-once message processing is kinda transactional
  • How exactly-once works
    • Approach 1
      • Wrap message ACK and database write in a single transaction
      • With support of distributed transactions, we can guarantee exactly-once semantics on different machines
      • This is possible ONLY if all systems affected in transaction support atomic commit - side-effects need to be reversed (example with sending email).
      • Example: XA
    • Approach 2
      • Only needs transactions in database.
      • Every message has its ID. When you database-processed a transaction, atomically save its ID to the table of processed IDs.
      • Send ACK to broker.
      • Once the message is ACKed by broker, you can delete it from your table.