Understanding LSM trees: What powers write-heavy databases (2020)
yetanotherdevblog.com
yetanotherdevblog.com
My unfinished attempt from 02014 to explain LSM-trees (and, in particular, full-text search engines using them) with an implementation in 250 lines of Python 2 is at https://github.com/kragen/500lines/tree/master/search-engine. I think it's a more complete explanation (except for tombstones), but it's longer (half an hour of reading) and there aren't as many illustrations.
On the plus side, you can actually run it:
$ make
./handaxeweb.lua index.py < README.md > index.py
$ python2 index.py index ../newindex .
$ python2 index.py grep ../newindex refactor
README.md:1173: # XXX refactor chunk reading? We open_gzipped in three places now.
index.py:220: # XXX refactor chunk reading? We open_gzipped in three places now.
$
I, uh, wouldn't recommend trying to run a high-performance service on my didactic database. It has both zlib and urllib.quote in the inner loop, and its merges are blocking.I think there are a couple of things in Groom's article that are not really part of LSM-trees per se. Like, bloom filters, or red-black trees, or caching sparse indices (skip files) in RAM. But that's true of mine too, which mixes the LSM-tree with the implementation of a full-text search engine!
The real killer feature is not write efficiency, but the sequential write pattern, as opposed to the small random writes used by B-trees. On flash, this makes a significant difference for endurance.
EDIT: And I forgot to mention that it is a lot easier to do compression on immutable data structures (SSTables) than on mutable ones, so LSMs are generally more space-efficient too.
Sequential writes with some degree of batching requests can add orders of magnitude to system throughput. I use a ringbuffer abstraction to serialize my writers and handle batches of up to 4096 transactions at a time. Everything in the batch gets compressed and written as 1 chunk to an append-only log. The benefits of processing multiple requests at a time are very hard to overstate. Saturating the write performance of virtually any storage device is trivial with this approach. It also carries back to magnetic storage really well, and even tape if you were so inclined to go down that path.
To get a better idea - Assume a 4K block size, and 128-byte size for some business type. You can hypothetically store 32 of these per physical write if you have perfect batching going on. Looking at a Samsung 960 which has ~2.1GB/s sequential write speed (which is approximately our use case), you would be able to persist ~16.4 million transactions per second in the most ideal scenario. This is extremely notable, because the maximum stated random write throughput for this device at QD32 w/ 4k block size is only ~440kops/s. Mix in even the most pedestrian of compression algorithms, and this 16m turns into an even more ridiculous figure assuming your transactions look somewhat similar to each other, the system is fully loaded, and the planets are appropriately aligned. These circumstances might sound extreme, but they are fairly common in areas like fintech where you might need to match millions of orders in ~1 second and its non-stop like this for hours.
I am not aware of how the LSM+WAL would meet the objectives I have laid out above - Most notably the implication that the disk is being touched multiple times per business transaction. Please correct me if I am mistaken in this regard. My solution ensures that the disk is touched <= 1 time per logical write (where the size of the write is bounded by the block size of the device).
This is not something specific LSMs implement, it's just how kernels do file IO.
RocksDB does also additional IO coalescing for concurrent writes, though that's more about reducing syscalls cost (one write() per write group, instead of one per write) than IO cost.
Ok - that makes sense. I think we are mostly on the same page here. My "SyncWAL" is invoked naturally as my buffer reader hits the barrier each time and dumps the current batch of items to be processed.
I can't think of any OLTP application where the writes are not random (in the sense of relative position in the key space). Either your primary keys are, or at least all your indexes are.
Write amplification is not an issue if sequential writes are orders of magnitude faster than random writes, which is usually the case, so even with that LSMs come out ahead.
Also if your application is read-heavy, you can tune the LSM parameters to have fewer levels (at the cost of more compaction traffic). With B-trees, you have to chase a fixed depth of log_{page size}.
LSM trees are a lot faster than B+trees with lots of random writes.
In a by-the-book B+tree implementation that doesn't allow MVCC, if your branching factor is 1024, your table is 256 gigabytes, and your records are 128 bytes, most random record writes require updating four tree nodes, the transaction log, and the actual record page, which would be six random seeks without a writeback cache, and probably three random seeks with a writeback cache. I think in practice on modern SSDs this means three 4-KiB blocks added to the FTL's transaction log, but without further write amplification I think that fragments the FTL's storage so that reading a tree node requires multiple random page reads instead of one. If you use a log-structured B-tree approach so MVCC is convenient, you probably want to use a lower branching factor; say you use 128. Now your tree is 5 levels of internal nodes deep, and each node is, say, 2 KiB of data, so you're appending 10 KiB of tree data to the write log for each random write, plus the 128-byte record.
By contrast, with an LSM-tree, at the moment of commit, you just append the 128-byte updated record, perhaps after an in-RAM sort: 80 times less data. But if the data survives long enough, it eventually has to get merged up to the root through a series of merges, which involves reading it and writing it an additional logarithmic number of times. Say your merges are 8-way on average; in that case the data has to get read and written 11 more times before the table is a single 256-megabyte segment. That's still 7 times faster than the log-structured B-tree approach.
The difference gets bigger with spinning rust.
And that's why, on one machine I tested on, LevelDB could insert 300'000 key-value pairs per second while Postgres manages about 3000 and SQLite about 17000.
For write amplification.
> If you use a log-structured B-tree approach so MVCC is convenient...
You just hamstrung your B-tree for conscience to implement a single model, but regardless: I never claimed that B-trees were faster at writes than LSM-trees. I instead claimed that I don't understand where all of these write-dominated use cases are coming from as the world seems to be full of read-dominated services.
> And that's why, on one machine I tested on, LevelDB could insert 300'000 key-value pairs per second while Postgres manages about 3000 and SQLite about 17000.
I mean, this is clearly an unfair comparison by your own math: you need to keep running this for a long enough time that you can amortize into your calculations the cost of compactions (for LevelDB) and vacuums (PostgreSQL).
What do you mean by "You just hamstrung your B-tree for conscience to implement a single model"? I'm not sure what you mean by either "for conscience" or "to implement a single model".
It's probably not an extremely fair comparison, but I'm not sure which databases it's slanted in favor of, and it wouldn't be very surprising if I'd overlooked some easy speedup. Like, maybe by deferring index creation until after ETL you could get a significant boost in Postgres speed, or by changing the Postgres tuning numbers. I do think it was running long enough to incur a significant number of compactions and vacuums, though. Feel free to critique in more detail. The Postgres and SQLite numbers come from http://canonical.org/~kragen/sw/dev3/exampledb.py, but for LevelDB I inferred that the Python bindings were the bottleneck, so the 300'000 figure is from http://canonical.org/~kragen/sw/dev3/leveldb-62m.cc.
Like we can assume the overwhelming majority of applications writing to a database are interested in durability. A good portion of those could tolerate eventual consistency on reads, for instance, which means there are considerably more ways outside the database to improve read performance.
For truly random data such as scale free graphs, B-tree shows quite non-linear performance decline, as if somehow it started to have O(N^2) operation cost instead of O(NlogN). LSM, on the other hand, worked very well there.
This approach can get you pretty far before you need sharding, auto- or otherwise, but the master's write bandwidth is the bottleneck. Taken to the extreme this realization leads you to ridiculous places like Kafka/Samza.
I'm trying to come up with an example of an OLTP application where the ratio of reads to writes within the same transaction needs to be large. I'm sure they must exist but I can't think of one.
It might make sense to take a different approach now than we did 20 years ago when we were mostly using MySQL and starting to adopt InnoDB. Caching reads and farming them out to readslaves is an easy performance win but probably introduces a fair number of bugs, and there's an enormous amount of powerful database software out there now that didn't exist at the time.
Note that even LSM-trees are only used for medium write intensity workloads. If you need maximum throughput for indexed access methods, newer database engines that use inductive indexing are on another level. I've seen a few implementations drive 10 million records per second through indexing and storage on an single EC2 instance (and with better query performance than LSM). These were not academic exercises either; they were invented to solve real-world problems LSM could not handle.
As for why anyone needs, say, 10 million inserts per second per server? There are quite a few modern workloads that generate many, many billions of new inserts per second if you run them at scale. Without the high throughput on a per server basis, managing the implied cluster size would become a serious challenge on its own. Server hardware can support very high ingest rates at very high storage densities if the software can drive it.
Can you provide me some resources where I can read more about this? Google is not helpful
Inserting an index into that pipeline means the data must be organized around a single primary index for every key column you want to query, no secondary indexes, and that you generally cannot use a classic multidimensional spatial index to capture multiple columns -- if column types have fundamentally different distributions of data and load, query selectivity degrades materially. A canonical example is trying to index time and space, and maybe an entity id, in the same index e.g. GPS probe data models, which can be both extremely large and extremely fast.
The general strategy is to build an indexing structure that adaptively compresses away all of the non-selective feature space of the dimensions regardless of the characteristic distribution of the data types. The idea is kind of old and follows from the theoretical equivalence between indexing and AI; whereas most succinct indexing structures effect compression via coding theory, this is compression via lossy pattern induction. Storage engines have to be specifically designed for this (ab)use case and you cannot do any kind of short-circuit queries with these indexes, you always have to look at the underlying data the index points to. Generalized induction is pathologically intractable, so these are efficient and practical approximations. The write (and query) throughput is extremely high because the index structure often fits entirely in CPU cache even for very large data models. AFAIK, these architectures have been used in production for around a decade, but current versions are qualitatively better than early versions.
That said, these indexes still suffer from the write sparsity problem that all cache-mediated index-organized storage has, it just pushes that performance problem out a couple orders of magnitude. Those limits are visible on ultra-dense storage available today. There is a radically different research database architecture that appears to qualitatively breach this limit but I don't expect that to show up in production software implementation for at least a few years. It isn't fully baked but people are working on it.
Virtually all of the research in this space happens completely outside of academia, so it rarely gets published, and research groups stopped filing patents over a decade ago because they were effectively unenforceable. It has turned into a boutique business.
If you’re thinking of putting a large write queue/buffer in front of the database to avoid losing data, well an LSM tree is one way to implement that. And since you can query it, you might not need a full database anymore. Read caches that lag behind the ground truth can do the job.
They make for both read- and space-efficiency.
In fact, that's a lot like what ZFS w/ ZIL is. The ZFS ZIL is not really an intent log, and its name is a misnomer. The ZFS ZIL is a log of past transactions, and the only intent is to merge them into the filesystem's B-tree-like structure.
B-trees are the best data structure for databases, except for write transactions: those are slow in a B-tree due to write magnification associated with the B-tree itself. A log amortizes the cost of writing to the B-tree. But the log also increases the cost of reads, so the log itself has to be indexed, which is what LSMs effectively achieve.
Wish it was just mentioned at least off hand that clearly this isn't the whole solution and a WAL or similar is used to make sure losing the memtable doesn't lose the data.
Assuming either that you only have one table, or file systems are magic things that don't make that argument moot.
> Once the compaction process has written a new segment for the input segments, the old segment files are deleted.
... and the assumption of only a single table isn't enough! Thank you, file system developers, for making SSTables sound so easy.
> Bloom filter
I did some research a couple of weeks ago about Approximate member query data structures. Bloom filters are from 1970, Quotient filters arrived in 1984, then Pagh improved Bloom filters in 2004, and we got Cuckoo filters in 2014 [1]. Has there been progress since? Did I miss any key achievement in-between?
Been meaning to work on some kind of LSM type hobby project in either Rust or Go.