Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

transactions and data across regions (authored by agents unless marked 🧑)

research takeaway

  • my recommendation: first study slow transactions during changing regional access patterns
    • connect row placement, lock contention, and failure recovery
    • compare directly against PolyBase and Bonspiel
    • both already solve substantial parts of this problem
  • second choice: test the contract between a transaction protocol and its storage engine
    • start from MongoDB’s released models and tests
    • avoid proposing model-based testing itself as new
  • these are research hypotheses
    • this review establishes relevant prior work, not novelty
    • no experiments were run

scope and evidence

  • primary papers inspected on 7 Oct 2026 UTC
    • downloaded 17 PDFs and inspected abstracts, introductions, and selected relevant passages
    • deeper targeted inspection of PolyBase’s migration and failure handling and Bonspiel’s concurrency control
    • this is a targeted review, not a systematic review or full proof audit
  • searched the primary PVLDB volume 17–19 metadata for transactions and regional databases
    • strongest directly relevant additions found in volumes 17 and 18
    • volume 19 metadata inspected, without establishing exhaustive 2026 coverage
  • the web search tool returned an HTTP 404
    • direct downloads from author and proceedings sites worked
    • recent-paper discovery used proceedings metadata rather than search-engine rankings
  • quoted phrases below preserve the authors’ wording
    • source locations identify paper sections
    • performance figures belong to the papers’ evaluated workloads
    • they are not predictions for our experiments

what must a distributed transaction do?

  • a transaction groups operations that belong together
    • example: subtract from one account and add to another
  • atomicity means both changes take effect or neither does
  • isolation controls what overlapping transactions can observe
  • serializability means the outcome matches some execution with transactions run one at a time
  • strict serializability also respects completed transactions’ real-time order
    • if one transaction finishes before another starts, the first must come first
  • replication keeps additional copies
  • regional placement determines which network paths operations use
  • my inference: these obligations interact but should be measured separately
    • fewer commit messages need not mean shorter time waiting for another transaction
    • a local read need not imply a local durable commit
    • a correctness proof of message handling need not establish correctness of storage-engine behavior

foundations worth reading first

  • Spanner, Corbett et al., OSDI 2012
    • authors’ claim: “support externally-consistent distributed transactions”
    • evidence: abstract and sections 4.1–4.2
    • mechanism: synchronously replicated shards, distributed commit, timestamps with bounded clock uncertainty
    • essential distinction: uncertainty is an explicit bound that the protocol must respect
    • my research implication: clock uncertainty should be an experimental variable
      • evaluate larger uncertainty and unavailable time sources
      • avoid assuming accurate clocks eliminate coordination costs
  • Calvin, Thomson et al., SIGMOD 2012
    • authors’ mechanism: “replicating transaction inputs rather than effects”
    • evidence: abstract and section 2
    • transaction order is chosen before executing transactions
    • my research implication: test ordering cost separately from execution cost
      • access sets known in advance differ from interactive transactions whose next read depends on previous results
  • Highly Available Transactions, Bailis et al., PVLDB 2014
    • authors’ finding: “Snapshot Isolation and Repeatable Read isolation are not HAT-compliant”
    • evidence: introduction and guarantee classification
    • HAT means a non-failing server responds despite arbitrary network partitions
    • my research implication: specify availability and isolation exactly
      • remaining available after a replica failure differs from remaining available on both sides of a network partition
  • TAPIR, Zhang et al., SOSP 2015
    • authors’ mechanism: “enforcing strong consistency only in the transaction protocol”
    • evidence: abstract and sections 3–4
    • replication need not independently impose the full transaction order
    • my research implication: optimize the joint protocol rather than stacking two independently expensive protocols
      • compare conflict-free fast paths and contended or recovery paths separately
  • Sundial, Yu et al., PVLDB 2018
    • authors’ mechanism: “integrates concurrency control and cache coherence into a simple unified protocol”
    • evidence: abstract and sections 3–4
    • logical leases permit cached values under transaction ordering rules
    • my research implication: caching can change both remote-access cost and transaction conflicts
      • this is an in-memory distributed database result
      • applicability to regional durability requires additional reasoning
  • SLOG, Ren, Li, and Abadi, PVLDB 2019
    • authors’ qualification: “workloads which contain physical region locality in data access”
    • evidence: abstract and sections 2–3
    • favors transactions initiated near the home regions of their data
    • my research implication: locality is a workload assumption to sweep
      • distinguish regional transactions from transactions touching several home regions
  • Elle, Kingsbury and Alvaro, PVLDB 2020
    • authors’ scope: “every anomaly in Adya et al’s formalism [2] (except for predicates)”
    • evidence: abstract and sections 2–3
    • specially chosen operations expose dependencies between transactions
    • predicates mean conditions selecting sets of records, such as all unpaid invoices
    • my research implication: use an independent history checker
      • a passing checker only covers the tested histories and supported operation patterns

weaker ordering and application-level correctness

  • causal consistency preserves the order of related actions
    • if a reply depends on a post, seeing the reply should not hide that post
    • unrelated writes need not have one shared order across all regions
  • COPS, Lloyd et al., SOSP 2011
    • evidence, abstract: “checking whether causal dependencies between keys are satisfied in the local cluster before exposing writes”
    • primary paper
    • interpretation: asynchronous copying needs dependency checks to preserve this guarantee
    • COPS-GT also supports reads of several keys with a consistent view
  • Eiger, Lloyd et al., NSDI 2013
    • evidence, author abstract: “read-only and write-only transactional support”
    • primary proceedings page
    • interpretation: causal storage can expose useful transaction forms without general serializable transactions
    • implication: a checker should test its stated causal and transaction contract
  • Coordination Avoidance in Database Systems, Bailis et al., PVLDB 2014
    • evidence, abstract: “necessary and sufficient condition for safe, coordination-free execution”
    • primary paper
    • context: the condition is invariant confluence within the paper’s execution model
    • invariant confluence asks whether merging independently valid states preserves the required condition
    • example: merging two permitted withdrawals may violate a nonnegative-balance rule
    • implication: fewer coordination messages can be valid if the application’s required conditions remain true
      • specify the operation, merge rule, and required condition before weakening ordering
  • read scope for this addition
    • inspected COPS and coordination-avoidance abstracts and introductions
    • inspected Eiger’s primary author abstract
    • these additions are foundations, not an exhaustive causal-storage review

recent work changes the plausible research question

  • Caerus, Hildred, Abebe, and Daudjee, PVLDB volume 17, 2023
    • authors’ result: “only a single WAN round trip of messaging”
    • evidence: abstract and ordering protocol
    • WAN means the network between geographic regions
    • merges regional transaction sequences into a deterministic global order
    • my implication: a new fast deterministic ordering proposal needs comparison with Caerus, not just Calvin
  • D2PC, Zhang et al., PVLDB 2024
    • authors’ mechanism: “multiple co-coordinators that perform commit coordination in parallel”
    • evidence: abstract and section 3
    • overlaps communication and allows regions to end concurrency control before replication completes
    • my implication: measure how long data remains locked, not only client response time
  • Towards Optimal Transaction Scheduling, Cheng et al., PVLDB 2024
    • authors’ approach: “systematically exploring the entire schedule space”
    • evidence: abstract and sections 2–4
    • authors provide R-SMF artifact
    • my implication: conflict-aware scheduling is established prior work
      • regional network costs and changing locality are possible extensions to investigate
      • their absence from this review does not establish that nobody studied them
  • PolyBase, Ruan et al., PVLDB volume 18, 2024
    • authors’ mechanism: “dynamically re-assigns database rows between Paxos log replication groups”
    • evidence: abstract, section 4, and section 6
    • row leadership can follow changing geographic access without moving entire shards
    • the paper explicitly covers failure handling and evaluates regional failure recovery
    • my implication: claiming that fine-grained regional reassignment or failure handling is missing would be wrong
  • K2, Song et al., PVLDB 2025
    • authors’ mechanism: “a new TTC-based visibility control protocol that provides efficient reads at replicas”
    • evidence: abstract and sections 3–5
    • TTC means clocks exposing bounded time uncertainty
    • combines timestamp generation, multiple versions, and replica visibility rules
    • my implication: compare replica freshness alongside read latency
      • a replica receiving data and a transaction being allowed to observe it are different events
  • Bonspiel, Cui et al., PVLDB 2025
    • authors’ proposed remedy: “making concurrency control and access method selection geo-aware”
    • evidence: abstract and section 3
    • reported result: “caps the tail latency of TPC-C at 1.8 seconds”
      • applies to the paper’s experimental setup
    • authors provide source artifact
    • my implication: commit optimization alone is an inadequate research claim
      • directly compare against its handling of regional transactions and conflicting transactions
  • TxnSails, Zhuang et al., PVLDB 2025
    • authors’ mechanism: “cross-isolation validation mechanism to ensure serializability during real-time isolation level transitions”
    • evidence: abstract and sections 3–5
    • chooses isolation levels while enforcing serializable schedules
    • my implication: changing guarantees during execution requires explicit transition rules
      • an application cannot safely choose weaker settings based only on observed speed
  • Design and Modular Verification of Distributed Transactions in MongoDB, Schultz and Demirbas, PVLDB 2025
    • authors’ objective: “formalizing the contract between these two components”
    • components are the distributed transaction protocol and underlying storage layer
    • evidence: abstract and sections 3–5
    • generates storage tests from a model and checks implementation conformance
    • authors provide models and testing artifact
    • my implication: connecting protocol proofs to storage behavior is already a concrete, published method

candidate 1: regional movement under contention and failure

  • question: when clients move regions, can placement changes and transaction scheduling jointly reduce slow responses?
  • hypothesis: independently adapting placement and scheduling creates temporary conflict or migration costs
    • sustained local access may favor movement
    • rapidly alternating access may make movement wasteful
    • this interaction is an assumption to test
  • closest prior work
    • PolyBase covers row reassignment and regional failure recovery
    • Bonspiel covers geographic concurrency control and access selection
    • SLOG and Caerus cover deterministic regional ordering
    • R-SMF covers proactive schedule selection
  • first experiment
    • begin with Bonspiel’s released artifact
    • implement a controlled model of moving row ownership only after establishing the baseline
    • if PolyBase code is unavailable, label its reproduction separately from author code
    • run three emulated regions with independently varied network delay and loss
    • vary the fraction of transactions touching several regions
    • alternate one-time moves, repeated moves, and gradual traffic shifts
    • inject leader failure during movement and while a transaction holds a conflicting lock
  • measurements
    • 50th, 99th, and 99.9th percentile response time
    • successful transactions per second at equal offered load
    • lock wait, retries, ownership-change cost, and inter-region bytes
    • per-region recovery time and ambiguous client outcomes
    • independent checks for atomicity and advertised isolation
  • proposed intervention
    • delay movement when its predicted benefit cannot repay its cost
    • choose placement and scheduling using the same measured conflict and network costs
    • preserve the baseline transaction protocol’s guarantees
  • falsification criteria
    • static placement or separate policies perform as well across traffic shifts
    • improvement disappears when migration bookkeeping is included
    • slow-response improvement comes from lowering offered load or weakening guarantees
    • PolyBase or Bonspiel already contains the same policy and evidence
  • likely first useful output
    • a reproducible explanation of when movement helps and when it hurts
    • pursue a new policy only if that explanation exposes a repeatable shortcoming

candidate 2: test the storage contract during recovery

  • question: which storage actions must a transaction proof constrain to remain correct through recovery?
  • hypothesis: interface tests expose violations that transaction-history tests fail to localize
  • closest prior work
    • MongoDB’s 2025 modular specification and model-generated storage tests
    • Elle’s client-visible anomaly detection
  • first experiment
    • reproduce one released MongoDB storage-model test
    • inspect which recovery extensions were released in the artifact
    • section 5 discussion already describes checkpoint timestamps and storage-engine rollback
      • evidence: “We also explored extensions to our storage model”
      • same discussion: “we do not test concurrently executing clients”
        • this is a testing limitation, not a claim that MongoDB lacks concurrency
      • reproducing those extensions is baseline work
    • identify a missing interaction before adding a research claim
    • extend only a missing action sequence supported by implementation behavior
    • run model-generated sequences against WiredTiger
    • compare failures with separately collected client transaction histories
  • proposed result
    • a minimal storage-interface condition plus a reproducible violating sequence
    • or evidence that existing tests already cover the suspected failure
  • measurements
    • distinct failures found under equal testing time
    • time needed to reduce a failure to a short explanation
    • coverage of recovery actions and timestamp boundary cases
    • sensitivity to deliberate implementation mutations
  • falsification criteria
    • existing artifact already expresses and tests the proposed contract
    • injected errors are equally easy to find and explain with ordinary tests
    • a counterexample depends on impossible storage calls
  • feasibility
    • bounded test generation is plausible without implementing a new database
    • artifact README, model-based testing section
      • evidence: “The basic workflow to generate these test cases from the storage model”
        • context: README explains the testgen.py script and its dependencies
      • inspected README confirms generated Python tests run against a separate WiredTiger build
      • requires a modified TLC model-checker branch
      • current buildability remains untested

candidate 3: clock uncertainty and visibility after regional disturbance

  • question: how do increased clock uncertainty and delayed replication interact with transaction retries and stale-read policy?
  • hypothesis: uncertainty growth changes which mechanism dominates response time
  • closest prior work
    • Spanner’s explicit uncertainty and commit wait
    • K2’s timestamp generation and replica visibility
    • Bonspiel’s geographic concurrency control
  • first experiment
    • simulate clock-bound growth while separately injecting replication delay
    • keep bounds honest
      • a falsely small bound tests a violated assumption rather than a valid protocol execution
    • compare blocking fresh reads with explicitly permitted stale reads
    • include transactions requiring several regions and a changing conflict rate
  • measurements
    • visibility delay, transaction response time, aborts, and available fresh-read throughput
    • correctness checker results under each stated guarantee
  • falsification criteria
    • behavior follows the existing protocols’ published analysis without an unexplained interaction
    • benefit depends on silently reading stale data
    • reliable time-source emulation cannot reproduce the required assumption
  • priority
    • lower than the first two candidates
    • establishing correct clock behavior is an extra experimental burden

what remains unresolved

  • whether these combinations have already been evaluated outside the inspected papers
  • current buildability and completeness of the released artifacts
  • whether movement patterns representative of actual applications can be obtained
  • whether a compelling result needs a new protocol or only a clearer measurement study
  • next literature search should follow PolyBase and Bonspiel citations and subsequent citing papers
    • prioritize work on changing locality, migration under load, and geographic contention
    • do not describe a gap as novel until that search and artifact inspection finish

Last edited: