And there is something amusing about using git as an example of how there is no technical argument against losing information. Because, losing information is a specific feature that was added to git. (Look up shallow clones.) I can accept that there is no technical reason to have this, actually. But there are pragmatic reasons to have it. Which is basically the point of this article.
This tweet from Rich provides the missing puzzle pieces:
"The log is not a btree - used for durability, not query. Separate indexes combine memory with batch-updated storage." https://twitter.com/richhickey/status/420910382538948608
Restating in my own words:
The transaction log is essentially just a stream of asserts and retracts with metadata. That's used for durability first and echoed to clients/peers and the indexing engine. The indexing engine asynchronously batch updates indexes (which probably are b-trees, but have no need to be streaming append-only style in this context). The peers query a joined dataset of the latest consistent indexes plus the set of unindexed transactions.
> As append-only B tree is persistent.
You're right, it can be. However, the popular append-only b-tree systems (such as CouchDB) are actually append-only B+-tree systems. That "+" means the leaves of the tree are intrusively linked together, and so are not persistent. You can not "fork" such trees cheaply.
> losing information is a specific feature that was added to git
Going further, you want fast queries on recent data. At this point, structures that build updateable indexes as fast as possible are going to be preferred. And you are likely looking at some form of b-tree for this. (Not necessarily, true. But likely.)
Excision is neat, but not necessarily the same thing. In git, if I do a shallow copy, I do not know that it was shallow. It literally gives me wrong information on who authored parts of code at this point.
And again, the point there is that it is a pragmatic feature, not necessarily a technical/ideal one.
Immediately. The approach used is to have a separate "long-term" data structure and a smaller short term one. The long-term structure is only batch updated periodically, while the short-term logs changes since last batch. Every access queries both.
And this doesn't even get in to the questions such as at what point is a record eligible to be in my query. When I ran my query, or when I iterated to where that record would be? Is this controllable?
And seriously, just answering "immediately" is very close to saying "by magic." Too close, for this old timer's preference. I have a ton of respect for the stack. More so for those making it.
The transaction log also goes into a tree, but that tree is structured rather differently (for performance reasons). For example, it maintains a "linked-list" (in storage!) of the latest N transactions, which it then rolls up into one tree node once that list gets to a certain size.
The missing part about how transactions are available "immediately" (which is a word that doesn't make sense for a distributed database ;), is that the transactor (which is a separate process/system from those that answer questions), streams new transactions (as they happen) to the query boxes (known as "peers"). A peer is just your usual client process: for example your Java frontend webserver process (at this time only JVM clients are properly supported in this model)
Indexing is done in the background every ~33mb of transaction data (in the transactor) (and it allows that to build up during new indexing jobs, applying back pressure if it gets too much in memory data). Indexing isn't append only at all - it creates a new tree (that very often shares a lot of data with the old tree, however).
To answer queries, the "peers" merge the new transactions they've received (that are in memory), with the durable index. That's how new transactions get seen quickly - the peers have recent data in memory, and other data in long term durable storage.
Transactions are visible as soon as the transactor's streaming sends them to peers. There is also a mechanism to say "wait until this transaction has arrived at this peer" before querying.
Because the data in the indexes is immutable, it's trivial to cache in the client processes. Many smaller databases can fit entirely in memory, in which case querying only hits main memory, not the network on the peer process, which makes them (potentially) many orders of magnitude faster than querying a traditional RDBMS).
I am curious on why you have "potentially" in parens. Is this just not panning out in measurements? Are these not techniques that older products could have already subsumed into their repertoire?
A lot of this reads like RISC versus CISC debates. There are virtually no techniques that one side can claim monopoly on. So it is not surprising to see that picking the acceptable tradeoffs and combining solutions appropriately is often very effective.
LMDB is a B+tree but it doesn't link leaf pages together.