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
- evidence: âThe basic workflow to generate these test cases from the storage modelâ
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: