Apple open-sources FoundationDB
foundationdb.org
foundationdb.org
The short version is that FDB is a massively scalable and fast transactional distributed database with some of the best testing and fault-tolerance on earth[1]. It’s in widespread production use at Apple and several other major companies.
But the really interesting part is that it provides an extremely efficient and low-level interface for any other system that needs to scalably store consistent state. At FoundationDB (the company) our initial push was to use this to write multiple different database frontends with different data models and query languages (a SQL database, a document database, etc.) which all stored their data in the same underlying system. A customer could then pick whichever one they wanted, or even pick a bunch of them and only have to worry about operating one distributed stateful thing.
But if anything, that’s too modest a vision! It’s trivial to implement the Zookeeper API on top of FoundationDB, so there’s another thing you don’t have to run. How about metadata storage for a distributed filesystem? Perfect use case. How about distributed task queues? Bring it on. How about replacing your Lucene/ElasticSearch index with something that actually scales and works? Great idea!
And this is why this move is actually genius for Apple too. There are a hundred such layers that could be written, SHOULD be written. But Apple is a focused company, and there’s no reason they should write them all themselves. Each one that the community produces, however, will help Apple to further leverage their investment in FoundationDB. It’s really smart.
I could talk about this system for ages, and am happy to answer questions in this thread. But for now, HUGE congratulations to the FDB team at Apple and HUGE thanks to the executives and other stakeholders who made this happen.
Now I’m going to go think about what layers I want to build…
[1] Yes, yes, we ran Jepsen on it ourselves and found no problems. In fact, our everyday testing was way more brutal than Jepsen, I gave a talk about it here: https://www.youtube.com/watch?v=4fFDFbi3toc
(I was one of the co-founders of FoundationDB-the-company and was the architect of the product for a long time. Now that it's open source, I can rejoin the community!)
i'm incredibly impatient to have a look at what the community is going to build on top of that very promising technology.
[talk] https://www.youtube.com/watch?v=4fFDFbi3toc
[by wwilson] https://news.ycombinator.com/item?id=16877401
In FDB's envisioned architecture, the "client" is usually a (stateless, higher layer) database node itself! So the client encompasses the first layer of distributed database technology, connects directly to services throughout the cluster, does reads directly (1xRTT happy path) from storage replicas, etc. It simulates read-your-writes ordering within a transaction, using a pretty complex data structure. It shares a lot of code with the rest of the database.
If you wanted, you could write a "FDB API service" over the client and connect to it with a thin client, reproducing the more conventional design (but you had better have a good async RPC system!)
thanks :)
The microservices crew with their "our database is behind a REST/Thrift/gRPC/FizzBuzzWhatnot microservice" pattern is still catching up to the significance of this statement.
First request the row needs to be read from disk HDD. It takes 2ms.
Second request, the row is already in ram, it takes microseconds but still has to wait for the first request to finish.
Threads have overhead when having a lot of concurrency (thousands/millions requests/second).
For extreme async, see seastar-framework and scylladb design.
TLDR: high concurrency, low overhead etc.
I don't think they will ever move 'secret sauce' into open source, but infrastructural things like DBs and dev tooling seems to be going in that direction.
Apple, like everyone else, wants to commoditize their complements.
> Apple chose to develop a new compiler front end from scratch, supporting C, Objective-C and C++. This "clang" project was open-sourced in July 2007.
Woolvalley said “and now recently clangd“, which is a source code completion server which began its life at Google.
Clang did start its life with Apple.
EDIT: My bad, looks like FoundationDB wasn't fully open-source back then.
You can even check the comments on the HN thread [1] when they removed the downloads.
I see no reason you wouldn't be able to implement Datastore. In fact here's a public source claiming that Firestore (which I believe is its successor) is implemented on top of Spanner: https://www.theregister.co.uk/2017/10/04/google_backs_up_fir...
Also TBH now that I don't have commercial reasons to push interop, if I write another document database on top of FDB, I doubt I'd make it Mongo compatible. That API is gnarly.
The Will Wilson I remember was not prone to such understatement. That API (and the corresponding semantics) was a nightmare.
Out of the many MongoDB criticisms, this one is valid. This https://www.linkedin.com/pulse/mongodb-frankenstein-monster-... article is quite right about it (note: endorsing the article does not mean I endorse its author by far).
Other than that, they totally did a fake it until you make it with MongoDB 3.4 passing Jepsen a year ago and MongoDB BI 2.0 containing their own SQL engine instead of wrapping PostgreSQL.
(I've tried to keep the above as dry as possible to avoid dragging the arguments around this situation into this thread - and I suspect the phrasing of the previous comment was also intended to try and avoid that, so let's see if we can keep it that way, please)
I realized this 2 months ago at https://news.ycombinator.com/item?id=16386129
Congrats to the MongoDB team!
Too funny!
It's fast to install on a dev machine though, lol.
It was a very bad rdbms (since they marketed as a replacement) and a very bad sharded db (marketed too).
FoundationDB is now licensed under Apache 2, which is a much more permissive license, so most companies' open source policies allow it.
I know that multiple other companies have a similar policy (either a complete ban on using AGPL-licensed software, or special approval required to use it), although unlike Google, they don't post their internal policy publicly.
If someone works at one of these companies, what do you want to do – spend your day trying to argue to get the AGPL ban changed, or a special exception for your project; or do you just go with the non-AGPL alternative and get on with coding?
It is true that MongoDB's AGPL contagion and compliance burden, if you don't modify it, is less than many fear. It is also true that those corporate concerns are valid. MongoDB does sell commercial licenses so that such companies can handle their unavoidable MongoDB needs, but they would tend to minimize that usage.
In any partition situation, if one of the partitions contains a majority of the coordinators then it will stay live, while minority partitions become unavailable.
Or do you get around this by not deploying an uneven number of boxes?
5 into 2, 2, and 1
7 into 3, 2, and 2
I think you can set the number of coordinators to be even, but you never should - the fault tolerance will be strictly better if you decrease it by one.
Anybody know if Apple migrated any projects from Cassandra to FoundationDB?
Or was every project using FoundationDB a greenfield project?
Would be great to see from Apple an engineering blog "strengths and challenges at scale" post now that it has been opened again.
The relative performance contribution of most of those features (cross-partition transactions not evaluated -it's a must) can be seen in our paper at Usenix Fast: https://www.usenix.org/system/files/conference/fast17/fast17...
FoundationDB uses optimistic concurrency, so "conflict ranges" rather than "locks". Each range is a (lexicographic) interval of one or more keys read or written by a transaction. The minimum granularity is a single key.
FoundationDB doesn't have a feature for indexing per se at all. Instead indexes are represented directly in the key/value store and kept consistent with the data through transactions. The scalability of this approach is great, because index queries never have to be broadcast to all nodes, they just go to where the relevant part of the index is stored.
FoundationDB delivers serializable isolation and external consistency for all transactions. There's nothing particularly special about transactions that are "cross-partition"; because of our approach to indexing and data model design generally we expect the vast majority of transactions to be in that category. So rather than make single-partition transactions fast and everything else very slow, we focused on making the general case as performant as possible.
Transaction coordination is pretty different in FoundationDB than in 2PC-based systems. The job of determining which conflict ranges intersect is done by a set of internal microservices called "resolvers", which partition up the keyspace totally independently of the way it is partitioned for data storage.
Please tell me if that leaves questions unresolved for you!
Ok, per my other question that makes sense. Similar to FaunaDB except the "resolvers" (transaction processors) are themselves partitioned within a "keyspace" (logical database) in FaunaDB for high availability and throughput. But FaunaDB transactions are also single-phase and we get huge performance benefits from it.
Systems that sound closest to FoundationDB's transaction model that i can think of are Omid (https://omid.incubator.apache.org/) and Phoenix (https://phoenix.apache.org/transactions.html). They both support MVCC transactions - but I think they have a single coordinator that gives out timestamps for transactions - like your "resolvers". The question is how your "resolvers" reach agreement - are they each responsible for a range (partition)? If transactions cross ranges, how do they reach agreement?
We have talked to many DB designers about including their DBs in HopsFS, but mostly it falls down on something or other. In our case, metadata is stored fully denormalized - all inodes in a FS path are separate rows in a table. In your case, you would fall down on secondary indexes - which are a must. Clustered PK indexes are not enough. For HopsFS/HDFS, there are so many ways in which inodes/blocks/replicas are accessed using different protocols (not just reading/writing files or listing directories, but also listing all blocks for a datanode when handling a block report). Having said that, it's a great DB for other use cases, and it's great that it's open-source.
I tried to explain distributed resolution elsewhere in the thread.
I believe our approach to indices pretty much totally dominates per-partition indexing. You can easily maintain the secondary indexes you describe; I don't understand your objection.
Also, the main draw-back of "indices as data" in NoSQL is when you need to add a new index -- suddenly, you have to scrub all your data and add it to the new index, using some manual walk-the-data function, and you have to make sure that all operations that take place while you're doing this are also aware of the new index and its possibly incomplete state.
Certainly not impossible to do, but it sometimes feels a little bit like "I wanted a house, but I got a pile of drywall and 2x4 framing studs."
This is a totally legitimate complaint about FoundationDB, which is designed specifically to be, well, a foundation rather than a house. If you try to live in just a foundation you are going to find it modestly inconvenient. (But try building a house on top of another house and you will really regret it!)
The solution is of course to use a higher level database or databases suitable to your needs which are built on FDB, and drop down to the key value level only for your stickiest problems.
Unfortunately Apple didn't release any such to the public so far. So I hope the community is up to the job of building a few :-)
Not even close. I don't even see anything I'd call a filesystem mentioned on your web page. I missed FAST this year, and apparently you had a paper about using Hops as a building block for a non-POSIX filesystem - i.e. not a filesystem in my and many others' opinion - but it's not clear whether it has ever even been used in production anywhere let alone become "best known" in that or any other domain. I know you're proud, perhaps you have reason to be, but please.
If you have a reason to use FoundationDB over HDFS, NFS, S3, etc, then this will work well.
Doing a Lucene+DB implementation where each entry posting lists are stored natively in the key-value system was explored for Lucene+Cassandra as (https://github.com/tjake/Solandra). It was horrifically slow, not because Cassandra was slow, but because posting lists are optimized and putting them in a generalized b-tree or LSM-tree variant will remove some locality and many of the possible optimizations.
I'm still holding out some hope for a hybrid implementation where posting list ranges are stored in a kv store.
FDB leaves room for a lot of creativity in optimizing higher layers. Transactions mean that you can use data structures with global invariants.
That sounds a lot like Datomic's "Storage Resource" approach, too! Would Datomic-on-FDB make sense, or is there a duplication of effort there?
Datomic’s single-writer system requires conditional put (CAS) for index and (transaction) log (trees) roots pointers (mutable writes), and eventual consistency for all other writes (immutable writes) [0].
I would go as far as saying a FoundationDB-specific Datomic may be able to drop its single-writer system due to FoundationDB’s external consistency and causality guarantees [1], drop its 64bit integer-based keys to take advantage of FoundationDB range reads [2], drop its memcached layer due to FoundationDB’s distributed caching [3], use FoundationsDB watches for transactor messaging and tx-report-queue function [4], use FoundationDB snapshot reads [5] for its immutable indexes trees nodes, and maybe more?
Datomic is a FoundationDB layer. It just doesn’t know yet.
[0] https://docs.datomic.com/on-prem/acid.html#how-it-works
[1] https://apple.github.io/foundationdb/developer-guide.html?hi...
[2] https://apple.github.io/foundationdb/developer-guide.html?hi...
[3] https://apple.github.io/foundationdb/features.html#distribut...
[4] https://docs.datomic.com/on-prem/clojure/index.html#datomic....
[5] https://apple.github.io/foundationdb/developer-guide.html?hi...
You can already tune segment sizes (a segment is a self-contained index over a subset of documents). I'd assume that the right thing to do for a first attempt is to use a Codec to write each term's entire posting list for that one segment to a single FDB key (doing similar things for the many auxiliary data structures). If it gets too big, then you should have tuned max segment size to be smaller. Do some sort of caching on the hot spots.
If anyone has any serious interest in trying this, my email is in my profile to discuss further.
I can confirm it wasn't fast!
(And to be fair that wasn't the point - back then there were no distributed versions of Solr available so the idea of this was to solve the reliability/failover issue).
I wouldn't use it on a production system now days.
[1] http://nicklothian.com/blog/2009/10/27/solr-cassandra-soland...
Out of curiosity, what led you to do this? And what does it do better/worse/differently than out-of-the-box things like Elasticsearch or SOLR?
I was in attendance at your talk, and thought it was one of the best at the conference. Apple I think broke some hearts completely going closed-source for a while, but glad to see them open sourcing a promising technology.
Kudos to the FoundationDB team and Apple for open sourcing this wonderful product. We're cheering you all along! And we look forward to contributing to the open source and community.
[1] https://www.snowflake.net/how-foundationdb-powers-snowflake-...
Do you have something to back that up? This to me reads like you imply that Elasticsearch does not work and scale.
It's definitely interesting but I'm cautious. The track record for FoundationDB and Apple has not been great here. IIRC they acquired the company and took the software offline leaving people in the rain?
Could this be like it happened with Cassandra at Facebook where they dropped the code and then more or less stopped contributing?
Also I haven't seen many contributions from Apple to open-source projects like Hadoop etc. in the past few years. Looking for "@apple.com" email addresses in the mailing lists doesn't yield a lot of results. I understand that this is a different team and that people might use different E-Mail addresses and so on.
In general I'm also always cautious (but open-minded) when there's lots of enthusiasm and there seems to be no downside. I'm sure FoundationDB has its dirty little secrets and it would be great to know what those are.
https://www.elastic.co/guide/en/elasticsearch/resiliency/cur...
https://aphyr.com/posts/317-call-me-maybe-elasticsearch
https://aphyr.com/posts/323-call-me-maybe-elasticsearch-1-5-...
The statement was blunt but the reputation is not exactly unearned. These problems become worse proportional to scale.
Unrelated to the original topic, but I had never come across that talk and it is great. I use the same basic approach to testing distributed systems (simulating all non-deterministic I/O operations) and that talk is a very good introduction to the principle.
When certain events happen rarely, you have fake them deterministically so that you have reproducible coverage.
To get great performance for SQL on FoundationDB you really want an asynchronous execution engine that can take full advantage of the ability to hide latency by submitting multiple queries to the K/V store in parallel. For example if you are doing a nested loop join of tables A and B you will be reading rows of B randomly based on the foreign keys in A, but you want to be requesting hundreds or thousands of them simultaneously, not one by one.
Even our SQL Layer (derived from Akiban) only got this mostly right - its execution engine was not originally designed to be asynchronous and we modified it to do pipelining which can get a good amount of parallelism but still leaves something on the table especially in small fast queries.
A tiny bit of caution for folks trying to run systems like this though: It is frigging hard at any reasonable scale. The whole thing might be documented / OSS and what not, but very soon you are going to run into deep enough problems that's going to require very core knowledge to debug, energy to deep dive. Both of which you probably don't want to invest your time into. Do evaluate the cloud offerings / supported offerings before spinning these up. Else ensure you have hired experts who can keep this going. They are great as a learning tool, pretty hard as an enterprise solution. I have seen the same issue a ton of times with a bunch of software (redis/kafka/cassandra/mongo...) by now. IMO In the stateful world, operating/running the damn thing is 85% of the work, 15% being active dev. (Stateless world is a little better, but still painful).
I remember all the people who bashed Apple when they acquired FoundationDB. I hope they are appropriately ashamed now.
I'm not ashamed about deriding apple for Apple taking a really, really great product and hiding it from the world for years to come.
This is definitely some atonement, but does not totally absolve Apple from the many times they've taken tech private.
The number of newbie engineers who see docker/kubernetes, and think "let me docker run or helm install" a stateful service in a couple of minutes - is mind boggling.
I'll remember this quote when trying to talk sense to them.
I'm probably a 1% engineer, been hired by M$, FB, and Google. These guys were light years ahead of me. I'm not sure I'm as good now as they were at like 17 years old. In fact I'm probably only a decent engineer from having observed the stuff they were doing back then and finding inspiration.
1: https://en.wikipedia.org/wiki/Thomas_Jefferson_High_School_f...
My second high school is unranked ("Texas Academy of Math and Science"). I don't think it qualifies as a high school. Seems we haven't done a great job of identifying the accomplishments of our alumni, based on the Wikipedia page. No doubt it would rank near the top though. My year alone Caltech accepted about 30 of us, more than any other high school in the country. Makes me wonder what my peers have been up to.
Anyway, I'd agree that these tech high schools have some amazingly smart people attending them.
2: https://en.wikipedia.org/wiki/Texas_Academy_of_Mathematics_a...
They aforementioned SWEs make themselves multimillionaires /and/ have great jobs /and/ get praised by their peers, yet now everyone else has to bow down to them... over claims they were great programmers in high school? Is an appeal to your own experience the best way to make yourself seem relevant here?
It's an incredible DB system, but it ain't the Second Coming. Calm down.
Is it ok to admire great engineers on HN?
Is this comment relevant to anything but your own inferiority complex?
Can we get a decent definition of hubris in here?
It's great to admire excellent engineers; aspiring to be as skilled as someone at a task can be very motivating. Worshipping them is another thing.
You're right about the inferiority complex - I know I'm a relatively bad SW engineer, but that's mostly related to how new it is to me. I expect and want to improve.
Hubris is defined as excessive pride or self-confidence. I would say that bragging that you're "probably a 1% engineer," and that you've been hired by three of the largest SW companies out there qualifies as hubris. Maybe it's just me, but I don't think a 1% engineer would publicly boast about being one and then actually use the 'M$' in a non-farcical manner.
Shitting on 99% of the SWE population to make yourself look good, then shitting on yourself to make another person look even better doesn't really work. There's a reason humanCmp() is a little more complex than strCmp().
BTW, being hired by a large company doesn't mean you're all that. Plenty of idiots get hired by Oracle.
1. It’s in its own repo
2. The build instructions are concise and clear. Dependencies are listed. You have to follow a total of 0 links.
3. They use a common build system and not an in-house thing.
Although it starts by running `make`, it's about as in-house as a thing can be.
https://github.com/apple/foundationdb/blob/master/Makefile#L...
And then there were the hijinks we went through to build a cross-compiler with modern gcc and ancient libc (plus the steps to make sure no dependency on later glibc symbols snuck in):
https://github.com/apple/foundationdb/blob/master/build/link...
Ahh... now that was a build system.
(with apologies to Tolkien)
Does the horse choke on the baseball? Is there an equine version of the Heimlich maneuver to be performed on horses suffering from mixed-metaphorical-adage-induced asphyxiation?
https://www.wavefront.com/wavefront-foundationdb-open-source...
Almost 5 years in and we have not lost any data (but we have lost machines, connectivity, seen kernel panics, EBS failures, SSD failures, etc., your usual day in AWS =p).
Or are you saying that AWS is particularly unreliable at scale?
https://cdn.chrisshort.net/How-Complex-Systems-Fail.pdf
Basically once a system is complex enough some part if it is always broken. The software must be designed from the assumption that the system is never running flawlessly.
This then causes people who aren't versed with the product or the technology to decrease their perception of the product, and puts the team behind it in a position of having to not just come to its defense but to do so quickly due to the perception concerns.
Meanwhile, if they just do a basic search for "MVCC serializable" they would find that they were wrong; which means that it took more time to leave this insulting comment than it would have taken them to learn how this can work.
As a community, we really really really really need to beat down on casual cynics like this, who like to lazily "call bullshit" or play the "citation needed" card as a way to undermine the credibilty of other peoples' products. We live in a future where the answers to these kinds of doubts are a moment away: this particular form of debate tactic needs to die.
I'll take your feedback under consideration, and I'm sorry to have irritated you.
A FDB transaction roughly works like this, from the client's perspective:
1. Ask the distributed database for an appropriate (externally consistent) read version for the transaction
2. Do reads from a consistent MVCC snapshot at that read version. No matter what other activity is happening you see an unchanging snapshot of the database. Keep track of what (ranges of) data you have read
3. Keep track of the writes you would like to do locally.
4. If you read something that you have written in the same transaction, use the write to satisfy the read, providing the illusion of ordering within the transaction
5. When and if you decide to commit the transaction, send the read version, a list of ranges read and writes that you would like to do to the distributed database.
6. The distributed database assigns a write version to the transaction and determines if, between the read and write versions, any other transaction wrote anything that this transaction read. If so there is a conflict and this transaction is aborted (the writes are simply not performed). If not then all the writes happen atomically.
7. When the transaction is sufficiently durable the database tells the client and the client can consider the transaction committed (from an external consistency standpoint)
The implementations of 1 and 6 are not trivial, of course :-)
So a sufficiently "slow client" doing a read write transaction in a database with lots of contention might wind up retrying its own transaction indefinitely, but it can't stop other readers or writers from making progress.
It's still the case that if you want great performance overall you want to minimize conflicts between transactions!
Is this similar to how Software Transactional Memory (STM) is implemented? It sounds very very similar indeed.
The basic approach isn't super hard to understand, though the details are tricky. The resolvers partition the keyspace; a write ordering is imposed on transactions and then the conflict ranges of each transaction are divided among the resolvers; each resolver returns whether each transaction conflicts and transactions are aborted if there are any conflicts.
(In general the resolution is sound, but not exact - it is possible for a transaction C to be aborted because it conflicts with another transaction B, but transaction B is also aborted because it conflicts with A (on another resolver), so C "could have" been committed. When Alec Grieser was an intern at FoundationDB he did some simulations showing that in horrible worst cases this inaccuracy could significantly hurt performance. But in practice I don't think there have been a lot of complaints about it.)
I don't know that there's been an exhaustive writeup of that part, but maybe one of us or somebody on the Apple team will put something together. It probably won't fit in an HN comment though!
Or... maybe this is the part where I point out that the product is now open-source, and invite you to read the (mostly very well commented) code. :-)
Going forward as an opensource product, I hope to see some clarity on the "how it works"... Distributed, performant ACID sounds good, almost too good to be true. Not that I doubt it at the moment, I just want to understand it better :)
MVCC needs a bit of additional logic ontop to be serializable - "Serializable snapshot isolation" is a good keyword to search for - But it's definitely possible.
https://courses.cs.washington.edu/courses/cse444/08au/544M/R...
https://drkp.net/papers/ssi-vldb12.pdf
https://wiki.postgresql.org/wiki/Serializable
Edit: Formatting
I believe that FoundationDB stores rows in lexicographical order by key. Other databases like Cassandra strongly push you toward not storing data this way as it can easily lead to hotspots in the cluster. How do you deploy a FoundationDB cluster without leading to hotspots, or perhaps what operational actions are available to rebalance data?
If you have a sorted data store, you can get the same distribution by keying off a hash of the "real" primary key, right?
Cassandra allows you to store multiple records in sorted order within a partition. The normal recommended way to get data locality is to store records that are frequently accessed together in the same partition.
I'm also kind of confused.. is the single repo complete?
What about the SQL Layer [0]? Where is all this stuff in the new GH repo?
Or is only the KV part be open-source?
Looking forward to some CockroachDB vs. FDB benchmark showdowns :)
[1] https://apple.github.io/foundationdb/known-limitations.html#...
"Wavefront by VMware’s Ongoing Commitment to the FoundationDB Open Source Project"
https://www.wavefront.com/wavefront-foundationdb-open-source...
https://apple.github.io/foundationdb/features.html
I wonder how it compares to MUMPS databases like Intersystems Cache and FIS GtM?
Also, the obvious difference is that Caché is closed source and prohibitively expensive.
FiS is only FOSS on selected platforms (Linux x86, OpenVMS Alpha), and proprietary on all other platforms.
Worked a job once where that was the underlying data store.
I was only allowed to touch the SQL interface to it, which was....weird.
The SQL dialect was ancient, felt like something from about 1990 (and this was in.... 2012 or so, so not THAT long ago).
Query performance seemed invariant. A simple select * from foo where id=X and a monster programatically generated join across 15 tables would both take about 1.5 seconds to return results.
Bloomberg’s comdb2 was open sourced recently https://github.com/bloomberg/comdb2 - it seems similar, but would be interesting to see comparison.
If not, then a lot of potential is being left on the table, because usage would require wrapping FoundationDB in a proxy or middleware of some kind to synthesize events, which can be extremely difficult to get right (due to race conditions, atomicity issues, etc). Without events, apps can find themselves polling or rolling their own pub/sub metaphor over and over again. If anyone with sway is reading this, events are very high on the priority list for me thanx!
I'm not sure it has every feature it will ever need in this area, but it's a pretty good starting point for building "reactive" stuff.
https://forums.foundationdb.org/t/log-abstraction-on-foundat...
from their source code blessing. notes.
Related snippet from the "Distinctive Features Of SQLite" page[1] from the sqlite project:
The source code files for other SQL database engines typically begin with a comment describing your legal rights to view and copy that file. The SQLite source code contains no license since it is not governed by copyright. Instead of a license, the SQLite source code offers a blessing:
May you do good and not evil May you find forgiveness for yourself and forgive others May you share freely, never taking more than you give.
Unfortunately it looks like they striped out some important things, notably the storage engine (there's now a sqlite fallback).
Edit: Apparently it was always sqlite as per replies bellow.
It's super easily pluggable[1], so now that it is open source people can experiment with other engines. I think there is a lot of room for improvement. Also architecturally it's designed in anticipation of being able to run different storage engines for different key ranges and for different replicas. For example, you might keep one replica in a btree on SSD (for random reads) and two on spinning disks in a log structured engine.
[1] https://github.com/apple/foundationdb/blob/master/fdbserver/...
It looks to me like Apple has made a pretty complete release of the key/value store. What's missing is
(1) Layers! Everything from relational databases to full text search engines to message queues
(2) Monitoring stuff. Unsurprisingly it doesn't look like we have the tools for monitoring log files, etc. Wavefront (also a major user!) is a great commercial solution, but there should be something OSS
Plus even more tooling (mostly Ansible) for managing large fleets.
Will you guys think about open sourcing tooling? Apple is realistically never going to do that stuff.
Can't speak to HBase, but one thing Cassandra doesn't guarantee is ACID - I've seen some data consistency issues that has arisen from Cassandra in our usage, although it hasn't been a huge problem for us. That difference alone probably brings a lot of value to FoundationDB.
I’d like to see a deep dive of how foundationdb handles this. It has to trade-off consistency for availability at some point and it would be nice to know exactly where.
From my memory writes within a row are atomic. It seems to pass Jepsen as well [1].
1. https://www.google.co.uk/amp/s/yokota.blog/2015/09/30/call-m...
Hbase 0.98 (I think?) introduced a feature called "timeline consistency" that allows reads from replica regionservers. This has to be enabled for a table and has to be specified on the query side. If it is, you have the option of falling back to the replica if the primary doesn't respond within a deadline. This may be a good tradeoff if you value availability over consistency.
The benefit of ACID transactions isn't just safety at the application level, it's the ability to compose abstractions and build complex data models on top of simple ones. For example, indexes in FoundationDB are typically more scalable than in the majority of other systems, where index queries often have to be broadcast to all the systems in a cluster. Yet FoundationDB doesn't even have indexes as a feature - higher layers build and maintain index invariants using transactions.
FoundationDB is also just really reliable and fault tolerant. Its testing story is drastically better than what the teams building these other products are doing.
It seems the announcement concerns only the KV one. Someone has information for the 2 other ones?
Thank you.
I didn't really pay much attention to Foundation before Apple bought them and am unsure how it fits in the wider database ecosystem.
But I would tell the story something like this: state storage is the root of (almost) all operational evil. It's very easy to make a system reliable if it's totally stateless. Even most bugs can be lived with if the worst you have to do is restart a service and carry on! But to do anything interesting you have to store state somewhere, and you have to modify that state concurrently without screwing it up.
And the many challenges of operating stateful systems are greatly multiplied if you have a lot of different ones. For example, if you have a datacenter outage and some but not all of your stateful systems deal with it correctly, probably your application as a whole is still down.
So as one more stateful system, does FoundationDB just make that worse? Well, FoundationDB is designed specifically to be a foundation for many very different stateful systems - not just different kinds of databases but things like search engines or message queues that you normally don't think of in the same category. So that almost any system can map to it efficiently, it has a lowest common denominator data model (key/value) and the highest possible guarantees in terms of concurrency control. And you can run diverse systems supporting an application on the same FoundationDB cluster, or on different clusters with the same exact operational requirements.
Some few users of FoundationDB have been able to get the benefits of this vision, consolidating lots of different stuff into a single, operationally desirable system. But for more people to be able to, not just does the key/value store have to be available to them, but also lots of stuff has to be built on top of it. By releasing FoundationDB under a very liberal open source license, Apple has hopefully made that possible. In the long run, hopefully it will make all server-side computing more reliable.
Also, it's a really good key/value store, if you happen to need one of those!
I noticed that all the write benchmarks in https://apple.github.io/foundationdb/benchmarking.html are for random writes. Is write throughput affected by highly-sequential writes (e.g. - time series) vs random writes? How do you avoid hot-spotting on recent ranges?
How efficient are range deletes?
On https://apple.github.io/foundationdb/performance.html I read "The memory engine is optimized for datasets that entirely fit in memory, with secondary storage used for durable writes but not reads." I'd like some clarification:
(1) Which memory does "entirely fit in memory" refer to? A single machine? Or SingleNodeMemory * Nodes / ReplicationFactor?
(2) If only recently-written data is likely to be queried, and all recently-written data fits entirely in memory, is that sufficient? If so, would an unexpected query of old data cause a huge impact on write throughput?
(3) What is the structure/format of the data stored on disk? How is it updated?
I'm wondering how well this could be used for time series data. I saw mention here that wavefront uses FoundationDB for this, but would like more details if any are available.
You can mitigate by designing your key structure/data ordering to not have that property.
The memory engine requires your data to fit in memory (total across all your nodes, after replication). It writes interleaved snapshots and updates to disk, and reads the whole dataset back into memory when restarted.
You can do great modeling of time series data in FDB, though it will take some care and thought.
You should ask these questions on the forum. This article is falling off HN, I am going to lose track of it, and it doesn't look like the Apple team is answering questions here.
Quick question: I know there's a watch API, but is there any way to subscribe to a change feed from foundationdb? I'd like to consume the FDB event log to do external indexing & map-reduce work.
Alternatively, your application or layer can use the "versionstamp" atomic operations to write its own ordered log of what it is doing, or other indexing tricks. Depending on your data model this might be able to be much more efficient. For example, for external indexing you probably don't need to preserve a history of prior values but only be able to identify all the values that have changed. This can be done with a very simple and compact index that doesn't need to duplicate all the data to be indexed.
I'm not sure I understand.
Are you suggesting having a second key space at `ops/{VERSIONSTAMP}` or something where values contain enough information about the operation to be able to process changes in an indexer? The indexer could then clean up after itself, deleting the operations once they had been ingested? ... Effectively using a portion of the keyspace as a queue?
If you are trying to make your external index MVCC, then you will want to carry some version information too.
This kind of question might be better served by the new community forum you can get to from the website!
In my understanding, FoundationDB's transaction management is closest to FaunaDB's; read/write sets are linearized in memory in preprocessing nodes and distributed asynchronously to the replicas rather than locked on the replica leaders like Spanner or CockroachDB. This is why FoundationDB doesn't support long-lived transactions.
It's interesting that the FoundationDB team chose to unwind their service architecture (there used to be separate transaction manager and replica processes), I assume in the interests of ease of operations.
It is not clear to me how leader election and failover works for the transaction management role. Maybe somebody from the team can clarify.
On the other hand, I don't really know.
FaunaDB and CockroachDB are implemented as monolithic processes and can break encapsulation boundaries for performance reasons. For example, FaunaDB does aggressive predicate pushdown to accelerate intersections and joins, which you cannot do if you have to conform to a key/value interface exclusively. It can also eliminate all network overhead for query data that's local to the processing node.
I understand how that terminology is confusing though...how would you explain it?
The CockroachDB documentation says their SQL implementation is layered on top of their K/V interface: https://www.cockroachlabs.com/docs/stable/architecture/overv...
This would make it similar to TiDB/TiKV and FoundationDB.
Here's and old article from VoltDB insisting that layering SQL on pure K/V forfeits too much performance: https://www.voltdb.com/blog/2015/04/01/foundationdbs-lesson-...
Edit: And a response: https://news.ycombinator.com/item?id=15505194
So new CloudKit codebase nowadays is being run on top of Foundation DB
Finally it's out! @WavefrontHQ managaes petabyte scale clusters with #foundationdb today!
ah found it: https://www.youtube.com/watch?v=oLGYMdo2q2g
But performance is going to suck if you run server nodes over unreliable connections. I have trouble seeing a FoundationDB cluster running on mobile robots as more than a trade show gimmick. Albeit an awesome gimmick. So in summary you should totally do that.
A couple of less time sensitive applications are: 1. distrubuting information the entire fleet should eventually know, 2. event log aggregation with fine-grained time alignment among nodes.
Both are probably silly problems to solve with a database, killing houseflys with sledghammers and all that, but it never hurts to explore creative tool misuse :)
The other option would be some sort or CRDT-based system - Antidote perhaps?
I'm positive that Citus & Spanner are quite different from FoundationDB, but I have no idea how. Googling didn't help much.
Can someone provide an overview of the differences?
Spanner (and to an extent its less mature OSS descendants Cockroach and TiKV) has more comparable goals, but is fairly different architecturally. For example, FoundationDB only requires N+1 replicas instead of 2N+1 to achieve N failure tolerance (even lots of databases with much weaker guarantees are in the latter category!), doesn't trust clocks at all, doesn't lose performance when transactions cross replica sets, and uses optimistic instead of pessimistic concurrency.
Also FoundationDB (and TiKV) make a distributed, transactional key/value store available as an API, while Spanner and Cockroach expose only a relational database layer. FoundationDB is designed philosophically with the idea that you want to have a single storage layer to manage operationally but should be able to mix and match data models and query engines above that layer.
On the other hand, FoundationDB doesn't currently have any full fledged high level database layer available. Someone will probably dig up our SQL layer (which was AGPL, I think) but I wouldn't really recommend using it in production because there is no active development team. Someone will probably try porting the SQL layers from TiDB and Cockroach.
Maybe Apple will open source more stuff in the future, but let's not get too greedy!
It can't tolerate N failures from N+1 nodes. It can tolerate N failures with N+1 copies of your data. In a big cluster you have plenty of nodes but storing everything 5 times to tolerate 2 failures is really expensive.
Sorry I got the terminology wrong, but that's a distinction without a difference. If it can tolerate N failures from N+1 copies, that means a network partition would allow any one copy to continue chugging along making changes by itself. You have two options: consistency is dropped and you downgrade to eventually consistent (at best), or availability is dropped meaning a single node can't make changes without a majority, which invalidates the N of N+1 failures claim. (Which is where the N failures of 2N+1 copies claim comes from in the first place: after N failures you still have a majority of copies.)
Or there's some other magic quorum protocol I've never heard of that makes the majority problem disappear.
In the happy case, replication takes place using the replicas and quorum rules specified by this configuration. For example, you might require writes to succeed synchronously against all N+1 replicas of some transaction log. After N failures, there will still be 1 replica remaining with the latest transactions. But in order to proceed after any failures, you have to do a consensus transaction against a majority of replicas of the coordination state, to specify the new set of N+1 replicas you will be using. And you also make sure that the 1 replica you are recovering from knows you are doing it, so that it won't continue to accept writes under the old replication configuration.
There can't be two partitions capable of committing transactions, because (in this case) you need either
(a) All N+1 replicas of the log, so that you can commit synchronously, or (b) A majority (N+1 out of 2N+1) of the replicas of the coordination state, AND 1 replica of the log
Sorry if this isn't a great explanation. Anyway it does work. I expect that you could rephrase this as an optimization of a consensus protocol, though I think it would be hard to build a performant and realistically featureful implementation that way.
When I wrote that I was wondering if it used a second 2N+1 dataset just for coordination & consensus. This has the benefit of separating data from consensus, allowing the N of N+1 data failure. But at the end of the day consistency still comes down to a N of 2N+1 failure tolerance of that second coordination state. It's smaller easier to replicate etc etc but it seems like it still has the same fault tolerance as just replicating the data 2N+1 times. It sounds like it's worked out great in practice for FDB.
But you say it rarely changes... but wouldn't it have to change every time there's a change to the dataset? I feel like this means you have to do even more replication and consensus than just replicating the data without this second consensus state.
In the best case with no failures this works great. But as the number of failures increases, I feel like due to the extra synchronization there will be an inflection point where the cost of the extra layers of coordination will be higher than just synchronizing the data directly. But due to 'other concerns' that inflection point is pushed back by a lot.
Is that a reasonable characterization?
It's certainly true that with enough failures you aren't going to make much progress. I'm not sure that is any less true with plain old state machine replication protocols, though.
Can you tell us what consensus algorithm you're using? Raft, or something else? Your own implementation or something off the shelf?
https://cormachogan.com/2014/04/01/vsan-part-21-what-is-a-wi...
The coordination consensus is (our own implementation of) disk paxos, which we liked for its operational properties in our context (the coordinators don't need to know about each other or communicate directly). An early version of fdb had a dependency on Zookeeper for this purpose; you can use anything.
Why is that?
[0] https://apple.github.io/foundationdb/configuration.html?data...
Most of the people who have run FoundationDB at scale have, for performance reasons, used configurations other than the "datacenter aware" mode for their inter region replication, so they may not be the strongest thing operationally.
There is some work that from what I can see in the code is still in progress to build a new, almost magical inter-region replication mode that I am very excited about, which combines synchronous replication to a "satellite" datacenter within region with asynchronous replication between regions and recovery logic that will finish replication and fail over in case of a partial failure of a region. You get fast transaction commits (much less than the inter region ping time), can fail over to a secondary region automatically and safely (without losing any committed transactions) in the vast majority of circumstances, and in the worst case you can (manually, because you are accepting data loss!) give up very recently committed transactions to fail over.
Thank you for your time and FoundationDB—along with @nlavezzo, and team(s)!
If a region is blown up instantly by an orbital laser cannon, then the database will go down and you will have to manually tell it to recover ACI in the other region, sacrificing the durability of whatever committed transactions in the lost region were destroyed by the laser cannon.
But what if you have different pieces of data and you want them to be fast in different datacenters? I think a great solution to this can be layered on top of multiple FoundationDB clusters, each using the satellite mode, but this is one thing that I at least haven't been able to think of a way to provide properly at the data model agnostic key/value store layer - the details about what to put where seem fundamentally dependent on your data model.
While local writes would stay fast, wouldn’t active/passive see higher-latency non-local writes than Spanner or Fauna’s (assuming a NAM-EUR-ASIA topology)?
I agree with and do appreciate the multiple FoundationDB clusters suggestion.
Wait, don't you need 3N+1 to tolerate N number of failures for it to be Byzantine fault tolerant? Is that not a goal of FoundationDB?
Citus I can't remember, but if Citus does it properly wrt consistency, it should be in the same category as CockroachDB and TiDB, there is no magic.
So if you actually put down your own fiber, and install GPS clocks in each rack, you’ll be able to enjoy the same results.
[0] https://www.cockroachlabs.com/blog/living-without-atomic-clo...
[1] https://cloud.google.com/spanner/docs/true-time-external-con...
[2] https://static.googleusercontent.com/media/research.google.c...
There’s still a lot to be said for the property, “scales horizontally to N machines,” of course, even if those machines have to be in the same data center.
From looking at the commit history it seems like this is pretty actively developed.
Seems dated as it requires NodeJS either 0.8 or 0.10.
Another step closer to a library appliance.
https://apple.github.io/foundationdb/performance.html#throug...
I wonder what the SSD engine performance would look with NVMe standard NAND or an Optane SSD instead of SATA. Any FoundationDB guys/gals on this thread able to comment?
Another Q: what's more commonly used in current FoundationDB deployments: memory engine or storage engine?
https://web.archive.org/web/20150304035646/http://blog.found...
Haven't scaled it yet to any large installations, so I can speak about how well it does that.
One oddity I see is that the ruby gem is not available on rubygems.org and therefore cannot be easily installed and maintained using the ruby package manager which is a bit of a pain.
I guess textbook SSI is willing to "reorder" conflicting transactions if the result is still serializable, which could violate external consistency if you don't have any other bounds on the order. In the language of SSI, fdb simply aborts the later of any pair of read/write transactions with an rw-conflict, in accordance with a fixed ordering which is externally consistent.
I guess it could also be that your book uses an idiosyncratic definition of linearizable, like trying to apply it to individual operations within transactions, which might rule out any optimistic concurrency method. It might just be better to delete this word from your vocabulary in the database field because there is no wide agreement on what it means. The first two hits on Google for me are Wikipedia and Peter Bailis, and they give clearly conflicting definitions, though I think fdb satisfies both!
Let me expand the definition in Kleppmann's book then.I think it is important because it creates a difference between SSI and typical Serializable level based on 2PL. The below is paraphrasing the definitions on p. 324-329. The book references http://cs.brown.edu/~mph/HerlihyW90/p463-herlihy.pdf. (I must admit, I read the book, not the paper).
Basic idea - make a system appear as if there were only one copy of the data and ALL operations on it are atomic. In this model, there may be replicas, but we don't care about them. As soon as a client completes a write to the db, all clients reading the db must be able to see the value just written.
In SSI this is not true, because you may the snapshot may not include writes more recent than the snapshot -> reads from the snapshot are not lineraizable.
Linearizable CAS register is equivalent to consensus, and can provide total order. It is therefore what most developers would love to have (if cost was not an issue :) )
"A history is serializable if it is equivalent to one in which transactions appear to execute sequentially, i.e., without interleaving... A history is strictly serializable if the transactions’ order in the sequential history is compatible with their precedence order... Linearizability can be viewed as a special case of strict serializability where transactions are restricted to consist of a single operation applied to a single object."
In these terms, FoundationDB has the strict serializability property, and thus if you do exactly one operation in each FoundationDB transaction then that is linearizable.
But that kind of linearizability is much less powerful than what FoundationDB actually gives you. You cannot efficiently maintain global invariants, like indexes, with single-operational linearizability. I don't think this definition is very useful! I think strict serializability (which is to say serializability & external consistency) is what you actually want.
A linearizable CAS register can be implemented in FDB as simply as this:
@fdb.transactional
def compare_and_set( tr, key, vold, vnew ):
if tr[key] == vold:
tr[key] = vnew
but this is not the limit of what you can do.for example this will commit under serializable:
create table counters(counter int);
insert into counters(counter) values(1);
BEGIN TRANSACTION ISOLATION LEVEL serializable;
select sum(counter) from counters;
/* insert sum into counters. wait until committing next transaction before executing the insert */
insert into counters(counter) values(1);
COMMIT;
/* this transaction should commit before doing the insert in the above transaction and after the above transaction has calculated the sum */
BEGIN TRANSACTION ISOLATION LEVEL serializable;
insert into counters(counter) values(10);
COMMIT;
both transactions commit and the final table looks like:1, 10, 1
which is possible if the first transaction committed first, and then the second transaction committed. but it is possible for another client to see the table as: [1], [1, 10], [1, 1, 10] which is a sequence of states which should not be possible. if you see [1], [1, 10] then you should see [1, 10, 11] as the last state. hence it violates external consistency.
Is that because PostgreSQL's "SERIALIZABLE" doesn't follow the "some serial order" definition? Or maybe I'm missing something else?
Serializability should ensure that the outcome is equivalent to some serial execution. What serial execution of those five transactions yields [1], [1, 10], [1, 1, 10]?
But also, I just read that prior to PostgreSQL 9.1 (released in 2011), the "serializable" isolation level was actually just snapshot isolation (now called "repeatable read"). So maybe that's what benmmurphy is referring to?
https://github.com/apple/foundationdb/blob/master/bindings/p...
Windows was never the most important deployment platform for the product. I'm not sure how performant the server is. But it should work, and I would think the client is fine.
You can also create graph layer over it using gremlin in a day or two.
Hope someone from the team reads this.
Emulating a low-level layer on a higher-level abstraction (which itself is using this hierarchy) will never match the speed, scale, or reliability of doing it correctly.
Interestingly, in the latest version of Ceph the abstractions are a bit different than you listed. Ceph is now using an object store built directly on top of raw devices. It's the file system, block, and object abstractions that exist on top of that.
https://github.com/spullara/nbd
Remember to format it with something like XFS rather than Ext to avoid writing superblocks all over the place.
I don't agree with this, but I think you may be confused because "Object Storage" can mean several different things.
"Object Store" in Ceph (as in RADOS - Reliable Autonomous Distributed Object Store) basically means key-value store. I typically say "blob store" instead to avoid the confusion with more sophisticated systems. It is exposed through a S3-like API. As far as I know, this layer of CEPH is pretty good, and you need a layer like this in most distributed systems anyway.
Ceph provides something called RBD, RADOS Block Device, which exposes a Block Device interface and is implemented on top of RADOS blob storage. It is useful for VM disks and has decent performance because it makes heavy use of the cache.
Some people use filesystems on top of RBD, but as far as I know CephFS itself does not sit on top of RBD. It is not as widely used as RBD because it is pretty recent (first release in 2016). The data is stored in RADOS and the metadata (which is the hardest part in a distributed filesystem) is dealt with by a Metadata Server cluster (MDS). This sounds like a typical distributed filesystem architecture to me, similar to GFS (the MDS replaces the GFS master and RADOS is used instead of chunk servers).
People tend to have a lot of issues with Ceph, but I think this is because:
1) It is used in reasonably large scale production settings where you are going to have issues anyway ;
2) It is not as easy to understand and fine-tune as it should be ;
3) Some people expect it to solve all their issues magically with perfect performance...
4) Some people use filesystems on top of RBD when they should have used CephFS or even direct interfaces to RADOS when possible.
But in general, I think Ceph is an example of a decently architectured complex distributed system.
Using Ceph for block and file access is like using AWS S3 to emulate block devices and filesystems. It'll work, and there is software for it, but it will never be very good. And Ceph is far from S3.
What are some examples of distributed file systems and block devices that _are_ very good?
That's all key to RBD being useful, or indeed CephFS itself. There are systems that map a filesystem layer on top of S3, but they have trouble because there aren't good ways to overwrite random small pieces of an S3 object. With RADOS, there are! :)
Yes, you can architect a storage system this way. But 1) Even if you do, many, many, many high-performance systems are "on top" of a filesystem but don't actually use the filesystem for anything except perhaps as a block allocator. Consider databases.
2) Many, many object stores do not abstract on top of filesystems. Modern RADOS, the distributed object store, stores its local data in a local object store called BlueStore. BlueStore speaks directly to the block device; there's no filesystem involved.
3) Even if you did store part of your distributed object store data on top of a local filesystem, that's not necessarily an issue. HDFS does this. (HDFS, despite the "FS" in its name, is an object store as most practitioners understand them.)
2) Rados is an object store, which is an abstraction on BlueStore (effectively a filesystem and replacement of FileStore), which is an abstraction on block devices.
3) HDFS is an object store, which is an abstraction on filesystems, which are an abstraction on block devices.
I'm not sure what you're point is because you just restated what I already said. They are abstraction, and they work just fine without any performance issues because that is the trade off of having an abstraction.
What I also said is that emulating low-level layers on a higher-level interface (like a block device on top of a database or object store) will never match the original block device. What is untrue about this?
One issue in this thread is that abstractions are concepts, not cpu instructions. In order to discuss overhead, one needs to reify the abstraction. For example, if you care about latency overhead, the block scheduler will definitely introduce overhead. But if you care about throughput, you probably /want/ abstractions like queues and schedulers.
> "What I also said is that emulating low-level layers on a higher-level interface (like a block device on top of a database or object store) will never match the original block device. What is untrue about this?"
Nothing is untrue about the sentiment of your statement. But from a practical standpoint, storage devices are useless pieces of junk without software. So to say abstractions slow down storage device while ignoring their utility feels arbitrary: why not talk about the length of the SATA cable, or the firmware in the disk controller? If the answer is that you just wanted to make the simple statement like the one I quoted at the start of this post then that's great, I think we are all in agreement. Otherwise, it's not clear what your point is and many of the supporting examples that you list are stated as fact, but are in reality either generally untrue, or very nuanced points, both of which tend to attract strong opinions :)
It certainly worked well within our required architecture - but as a general purpose system it would have several issues over typical network topologies.
As active volume size increases and/or changes, the object storage layer latency can bump up to minutes, or even hours. We stopped tested at 100TB volumes, as the object storage layer was backed up to hours.
Object Storage is very convenient, but the lack of good metadata and latencies involved basically means that it's only good as an archival backing store of active data. If your active datasource goes out - you could lose hours of data unless you have local copies.
The speed and durability has gone unmatched. We can write millions of transactions per second w/ millions of reads w/o issue. We have never lost data or found data to be inconsistent.
I make no comment on the validity of the stance but I think you are probably in the minority.
However since it is now open source at least it is possible for someone to do the work to get it on FreeBSD at least.
Maybe these two would go good together? ;)