The end of a myth: Distributed transactions can scale
muratbuffalo.blogspot.com
muratbuffalo.blogspot.com
> The data is assumed to be randomly distributed to memory nodes in the shared memory pool. Any memory node is equi-distant to any compute node, and a compute node needs to reach multiple memory nodes for transaction execution.
and also:
> Once a memory server fails, NAM-DB halts the complete system and recover all memory servers to a consistent state from the last persisted checkpoint. The recovery procedure is executed by one dedicated compute server that replays the merged log for all memory servers.
>This global stall is problematic for production, but NAM is a research proof-of-concept system. Maybe using redundancy/quorums would be a way to solve this, but then that could be introducing challenges for consistency. Not straightforward.
Yeah, research and proof-of-concept sounds about right. Some might even call it a toy database.
Also worth noting that this is a paper from 2016/2017. The world of (actually!) distributed database design has significantly moved on since then.
What are some of the significant movements since then?
distributed transactions that can assume fast and reliable access to some shared memory are interesting but really the easy part of the problem
"scale" means going across large physical distances where neither low latency nor high availability can be assumed, so not in a single datacenter where nodes are connected by infiniband or whatever
Way back in 2015, MySQL Cluster (NDB Cluster engine) benchmarked 200m transactions/second on commodity hardware [1]. It was read-committed transactions, not snapshot isolation, but still impressive. NDB (or RonDB, the new DB by its author) uses a non-blocking 2-phase commit protocol (failed transaction coordinators are failed over) and is even open-source. Still, it hasn't had a big impact outside of the Telecom and real-time gaming worlds.
[1] https://www.slideshare.net/frazerClement/200-million-qps-on-...
FoundationDB, TiKV, CitusData's Postgres extension are very well known, but RonDB looks like a hidden gem.
What are the practical issues with RonDB that it is not widely known(or used?)?
I think reasons for slow adoption are probably a mix of:
1. Lack of developer awareness.
2. Security implications (or perceived implications) of exposing memory directly to a network without passing through CPU or application-level access control mechanisms.
The second point may be prohibitive for a lot of general purpose database systems which are intended to be run on shared infrastructure on virtualized instances.
Another reason may be that a lot of production systems are CPU-bound, not memory-bound. RDMA seems ideal for systems which require a lot of memory. I'm thinking maybe with recent advancements in AI/LLMs, it could be an interesting technology as these do require a huge amount of memory relative to CPU.
The experimental setup involves using a cluster of 56 machines connected by an InfiniBand FDR network
That bit might have something to do with it.A bulk purchase of ~60 FDR IB cards, cabling, and network switches to support them sounds pretty expensive.
That being said, IB FDR gear is "older tech" now so the cards and switches can commonly be found reasonably cheaply on Ebay. The switches tend to be bloody loud though, so they're not something you'd want nearby if you can help it.
Basically, we are stuck in a local minima of making databases work well in cloud environments that are not designed to enable efficient distributed databases. It is wasteful and also provides an arbitrage opportunity for cloud companies.
- this paper
- dedicated bandwidth (Azure gives you bandwidth based on instance size)
- XDP on network interface
Probably more, but those are what I know of from running into their non-existence.
So while the spanner paper is open to be implemented by any vendor, they don't have the proprietary advantage that Google has - the atomic clock. So Yugabyte, CockroachDB don't rely on atomic clocks. I tried to get to the ground level basics of this, but I haven't understood this matter completely yet.
Note that even Spanner had multiple downtime due to clock and/or network failures. In these case, any operation lal guarantees are lost. This makes it really dangerous.
Systems limited by memory as in "quantity of" are scarse
The idea that one of many writer-compute-nodes can literally reach into a memory buffer that is shared across machines, atomically flip some lock bits and propagate some cache-coherence messages, and use that to build a multi-writer distributed database without needing to partition (and where any writer-compute-node can handle any message, so you can just round-robin a firehose of messages at them)... and that there's a chance (though not yet implemented) that one could implement ACID on top of this? It's absolute madness, and wildly exciting.
At least now I can provision Cloud Spanner as a managed service, but is this the future of clouds keeping their services at an advantage?
It’s amazing what you can get, if you just ask. It’s the latter part you really want, which is doable with modern hardware accelerated time stamping and a peering NTP system (like chrony).
Sarcasm aside, their NTP servers are less than 2 hops from my router. I know that because of excellent tools like ping and trace route.
With a correctly equipped and configured PTP installation locally, you can expect your clocks to be synchronized on the order of a very small number of microseconds with respect to each other, but the relationship between those and the atomic global clock is something else.
If you are using satellite/GPS time sources, then that's another matter. Why then is it relevant to have an NTP hookup at your ISP?
Also, consistency and accuracy are orthogonal to each other. My servers clocks are consistent and fairly accurate in regards to global time.
crazy customer needs is what made up most projects I saw or worked with.
Most machines will also know it is 8:11.
The question is how much resolution do you need and are your clocks accuratly synced enough for the thing you plan to do.
You are aware that many of the physics experiments that operate at the edge of what is possible use extremely accurate clocks synced over national and sometimes continental borders?
¹ let's ignore the ones that have been configured incorrectly
Fairly certain distributed physics experiments will be using atomic clocks that are not synchronized over the network. And/or they'll use direct satellite or GPS time sources which have predictable latencies.
You can't trust system clocks in a distributed system to ensure ordering. Some reading:
https://codeburst.io/why-shouldnt-you-trust-system-clocks-72...
My understanding of relativity is that this is principally unknowable and we only assume that speed of light is same in every direction by convention (every few years there’s a paper that claims they managed to measure one way sol but later it turns out they actually measured two-way in a roundabout way) so you can’t really know that?
PTP has good support in higher-end switches and routers, but it's difficult to secure and make resilient to failures. It was designed for automation and control networks in factories etc. NTP is a better fit for computer networks, but there doesn't seem to be any switches or routers with HW NTP support. If you really need the best accuracy with NTP, you can find old 100Mb/s hubs on ebay and create a separate network.
There's no networking hardware timestamp support for NTP because NTP has nothing to do with hardware timestamps.
PTP can be done without hardware timestamps, but it was designed with hardware support in mind.
I don't know where you got it that NTP does anything even orders of magnitude close to nanoseconds:
> NTP can usually maintain time to within tens of milliseconds over the public Internet, and can achieve better than one millisecond accuracy in local area networks under ideal conditions
https://en.m.wikipedia.org/wiki/Network_Time_Protocol
> The Precision Time Protocol (PTP) is a protocol used to synchronize clocks throughout a computer network. On a local area network, it achieves clock accuracy in the sub-microsecond range, making it suitable for measurement and control systems.
https://en.m.wikipedia.org/wiki/Precision_Time_Protocol
Literally neither solution comes anywhere near nanosecond accuracy.
For reference, there are 1,000,000 nanoseconds in a millisecond
Both NTP and PTP don't care (as protocols) where the timestamps are coming from. That's an implementation detail.
> NTP can usually maintain time to within tens of milliseconds over the public Internet, and can achieve better than one millisecond accuracy in local area networks under ideal conditions
That was maybe 20-30 years ago, but not today. The wikipedia article needs an update. If you don't hit a routing asymmetry, in my experience it's usually milliseconds over Internet and tens of microseconds in local network if using SW timestamping. Please note that NTP clients by default use long polling intervals to avoid excessive load on public servers on Internet, so they need to be specifically configured for better performance in local networks.
You can find some measurements with HW timestamping here: https://chrony.tuxfamily.org/examples.html
Note that this is for the system clock, which has to be synchronized over PCIe to the hardware clock of the NIC. That adds hundreds of nanoseconds of uncertainty. It doesn't matter if the hardware clock is synchronized by PTP or NTP.
If you care only about the hardware clocks, it's easy to show how accurate is the synchronization by comparing their PPS signals on a scope. NTP between two directly connected NICs, or a with a hub, can get to single-digit nanosecond accuracy. I have seen that in my testing. It's just timestamps, it doesn't matter how they are exchanged.
PTP is often used now in the Telco world as well as for broadcasting applications.
I in fact think once you offer atomic clock as a service for any public cloud, distributed transactions become a lot easier to implement thereby disrupting Spanner.
It reminds me specifically of 90's multi-client LAN database systems (dBase, Clipper) where clients coordinated via file locks. Unreliability & hangs became a big problem for us.
In the summarized RDMA database, I'd be pretty concerned about reliability & integrity:
1) Crashed servers will leave records locked, and the system will hang. 2) Question whether lock timeouts can be adjudicated reliably. 3) Any errors in server behaviour can easily & widely corrupt data across any other nodes. 4) Overall the RDMA coordination makes me cautious. Can we really replace Paxos with RDMA reliably? If not, problems squeeze out elsewhere. 5) Proposed single-threaded recovery procedure sounds a hazardous operational bottleneck. 6) I'm also cautious about coordination requirements around recovery/ or to transact knowing that recovery is not in process, unless we can show that can be reliable & not add cost to the protocol.
I implemented multiversion concurrency control in Java with threads but I want to raise it to multimachine.
The problem I have is that the read and write timestamps need to be available to detect if a transaction between machines conflicts.
How do I synchronize read and write timestamps with minimal latency?
When a database replica receives a transaction for key X, it needs (a) a timestamp that is globally accurate (b) needs to tell other replicas about it so they can detect dangerous read-write dependencies.
I feel it inherently requires serialisation of timestamps to accurately determine what other nodes are doing with data.
Otherwise you get dangerous read-write skew.
(See the whitepaper "Serializable Snapshot Isolation")
Load balancers don't do much work, except shift traffic. I wonder if a load balancer that collects timestamps and enriches requests with timestamp information would be a scalable solution?
Also see https://tikv.org/deep-dive/distributed-transaction/timestamp...
Cockroachdb describes their approach here: https://www.cockroachlabs.com/blog/living-without-atomic-clo...
Everything is direct memory access.
SNA protocol, No TCP. This is used since the 70s by many banks to process high volume transactions
RDMA (remote direct memory access) is a zero-copy communication standard.
RDMA is a user-space networking solution, accessed via queue pairs: lock-free data structures shared between user code and the network controller (NIC), consisting of a send queue and a receive queue.
RDMA supports several modes of operation [such as] reliable two-sided RDMA operations, which behave similarly to TCP. With this mode, the sender and receiver bind their respective queue pairs together, creating a session fully implemented by the NIC endpoints.
Once a send and the matching receive are posted, the data is copied directly from the sender’s memory to the receiver’s designated location, reliably and at the full rate the hardware can support.
A completion queue reports outcomes. End-to-end software resending or acknowledgments are not needed: either the hardware delivers the correct data (in FIFO order) and reports success, or the connection breaks.
> Is RDMA mature (robust/reliable) enough to use in distributed transactions?
At least on Linux systems, RDMA has been pretty robust/reliable for probably a decade, maybe more. > What are the handicaps?
It works differently to TCP/IP, which everyone in IT has at least passing familiarity with. So, it tends to be automatically passed over unless people hit a situation where they're open to "exotic" solutions.That being said, there's a TCP/IP shim layer (IPoIB) available which can be used by existing software to run on an IB network.
That shim layer though used to have a reputation for flakiness, and it was (or at least used to be) measurably slower than using native IB.
> What are the reasons for slow uptake on this?
Repeated self-inflicted foot-guns by Mellanox leadership or perhaps their sales and marketing leadership is my best guess.Mellanox adapters when brand new are priced fairly high for network adapters, or at least they used to be. However, when a network using them is upgraded to the next generation gear, a substantial number of the old ones would commonly end up for ~cheap resale on places like Ebay.
So, *nix DevOps staff ("Sysadmins" back in the day) and anyone else that needed fast networking for their home labs and similar would pick them up and figure out how use them.
Which of course meant over time an increasing group of people familiar with IB, that when figuring out solutions for their work places would then have enough confidence and knowledge to order them.
Sounds like typical organic growth right?
<rant>
Except Mellanox leadership - or at least their Sales & Marketing people - seemed to be fucking horrified that people were buying their expensive adapters cheaply on Ebay.
So, while the Mellanox technical people were receptive to this ground swell growth in usage by "unofficial" people, and tried to help out, their leadership time and time again did their level best to stamp it out.
Including killing off one of their decent "Community" growth initiatives -> just turned it into a fucking marketing channel for press releases and other crap <- Ugh.
But also doing things like instructing their support staff to not answer anyone on their forums who seemed to have purchased their gear through unofficial channels, etc. Nor let anyone else do so.
There were many, many examples of this bullshit over the years. And they were always like "Why don't we have huge adoption?"
Gee, I wonder? :( :( :(
</rant>
Anyway, I gave up on them and moved on a few years afterwards, prior to Nvidia buying them. I still buy older Mellanox ConnectX-3 (VPI) adapters off Ebay occasionally for home gear use (with Linux), and they're still solid. 10/40GbE. :)
The two papers mentioned in the post are highly insightful. This is one of my favorite subjects.
To fully disclose, I have a biased view on this given that I work for a (closed source) serverless, no-ops DB provider (Fauna) that implements a distributed transaction engine that is natively document-relational and doesn't compromise on relational (ACID, transactional) guarantees.
The particular timestamp oracle implementations still might suffer from metadata bloating and its commit protocol assumes an optimistic segment lock which would be a bottleneck in multiple scenarios.
So, from what I understand this is a (good?) practical speedup for some scenarios but doesn't change the fact that it's not possible to make an arbitrary sytem distributed without sacrificing correctness or throughput.
An easy reasoning: imagine that every node is one light year apart from all others.
I’m curious to see what advantages this paper offers to everyday Cloud Developers over using the latest version of Cassandra.