11 ms·
How FoundationDB works and why it works (2021)
- alberth 3y agoGreat article. Demystified a lot about FDB for me. > ”Summary: FDB is probably the best k/v store for regional deployment out there.” Why should someone use Memcache or Redis then? Is it for the data types in Redis?
- c0balt 3y agoRedis has a few features outside of k/v, like a good pub-sub implementation, that make it very useful in addition to a good DX and mature libraries. Memcache on the other habd is just solid and mature. It also has some inertia as being a solid k/v cache. For example: NextCloud supports afaik both Redis and Memcache as caching engines but doesn't have FDB support.
- deleted 3y ago[deleted]
- rapsey 3y agoBecause memcache and redis are in-memory. Writing to fdb will be complete once it is fsync'ed to disk. Memcache is a cache. Fdb is a an ordered kv store.
- abhishekjha 3y agoYou can configure redis to flush to disk on write operations though you lose on performance.
- leetrout 3y agoThat is not comparable to writes completing after they are flushed to disk.
- NavinF 3y ago>appendfsync always: fsync every time new commands are appended to the AOF. Very very slow, very safe. Note that the commands are appended to the AOF after a batch of commands from multiple clients or a pipeline are executed, so it means a single write and a single fsync (before sending the replies). https://redis.io/docs/management/persistence/ https://redis.io/docs/management/persistence/ It's very slow, but if you really want to wait for fsync before replying, it can do that.
- leetrout 3y agoI was unaware they could make that guarantee. Thanks for the correction.
- unmole 3y agoRedis can be configured to persist and fsync every operation.
- rapsey 3y agoRedis is not meant as a primary database and should not be used as such. FoundationDB is meant as a reliable source of truth.
- Already__Taken 3y agoRedis began as a caching database, but it has since evolved into a primary database. Many applications built today use Redis as a primary database. https://redis.com/blog/redis-cache-vs-redis-primary-database-in-90-seconds/ https://redis.com/blog/redis-cache-vs-redis-primary-database... Very much seems like an acceptable use now
- vore 3y agoRedis can't have a working set larger than memory. It has no mechanism to page data to disk. If your data set grows too large, you're hosed unless you add more hardware.
- dboreham 3y agoThings like Redis and Memcache are not serious data stores. Don't put any data in them that you really need back out later.
- rullopat 3y agoFDB is more a framework to create your own distributed database creating what they call a "layer". https://apple.github.io/foundationdb/layer-concept.html https://apple.github.io/foundationdb/layer-concept.html
- leentee 3y agoI believe the author means "the best transactional k/v store"
- abhishekjha 3y agoThis is the second article after the "caddy" one that I am having troble finding a usecase. Nginx eixsts, why do I need to learn caddy? Redis exists, why do I need to learn FundationDB?
- deleted 3y ago[deleted]
- qaq 3y agoIf you don't have a use case that requires FoundationDB you def. do not have to learn it.
- abhishekjha 3y agoWell, for that the comparative features needs to be clearly laid out. It solves a problem that others have not solved. Which problem is that? And how better does it work than whatever there was?
- mst 3y agoIt's an analysis of a paper, not a feature checklist, so the article is for people who want to understand more about how FoundationDB works rather than trying to provide a "which system should I use" explanation. Though "super reliable distributed database presenting a single logical shard to client code" is a class of system for which I don't think there's anything else even close out there, and I suspect generally if that's something you -need- then you'll already know that.
- qaq 3y agoIt's a highly scalable (as in powers iCloud services with billion users) distributed transactional KV store. It's owned by Apple and mostly developed by Apple and Snowflake.
- rapsey 3y agoFoundationDB replaces MySQL/PostgreSQL (if the tradeoffs are acceptable) or Cassandra. It is a reliable distributed store. Unless you are running Redis only with nothing else, fdb and redis do not play in the same space.
- jrvarela56 3y agoTangent about FoundationDB: this is a great video that explains how the team tested it https://youtu.be/4fFDFbi3toc?si=kSZ8VcOIjW_pMmPd https://youtu.be/4fFDFbi3toc?si=kSZ8VcOIjW_pMmPd Spoiler: they even had their own custom power supplies used to test against power failures.
- dboreham 3y agoFwiw this is standard procedure for anyone shipping a persistent storage product.
- rapsey 3y agoOnly after FoundationDB made it standard.
- mst 3y agoFoundationDB's testing turns the rigorous up to 11 and I'm unaware of anybody else who's published a description of a testing approach that goes to quite the same extremes. If that's just because I haven't noticed the others, I'd love to hear about them for comparison.
- jeffbee 3y agoI found a dataloss bug in the first hour of testing FDB. I think their testing hype is a bit overhyped. I also find it somewhat irritating that they won't take fixes or reports of problems with the storage engine because "we fixed this in redwood" when redwood is completely theoretical.
- endisneigh 3y agoYou've been able to use redwood since FDB 7.1 via ssd-redwood-1-experimental. Hardly theoretical. Curious about that data loss bug. do you have a link? most bugs I've seen have to due with latency spikes and cluster unavailability. haven't seen any around data loss after transaction has committed.
- deleted 3y ago
- jingles_dev 3y agowhen do we choose FDB over Redis and vice versa?
- datadeft 3y agoFDB is for permanent data storage, Redis is for temporary data. It does not matter that you can configure Redis to persist data to disk because the performance most Redis use cases need makes it less usable that way.
- webmonkeyuk 3y agoI suspect: - memcached if you don't need to persist the data - Redis if you don't know whether you need to use Redis or FoundationDB - FoundationDB if you learn that Redis doesn't do what you need I don't mean this in any kind of a derogatory way but I suspect that if you need to ask then you probably don't need FDB. The principle of keeping tech stacks boring and using well established components is less exciting as an Engineer but is usually the best choice.
- jamesblonde 3y agoIf you need transactions and high availability, then use FDB. But if you also need low latency / high throughput, then you should consider RonDB Disclaimer: i am involved with RonDB
- LAC-Tech 3y ago"FDB can tolerate f failures with only f+1 replicas." Wait a minute, I know that formula... Viewstamped Replication?? I need to read the foundation DB paper. (I mainly read CRDT stuff so hopefully it's understandable). --- In general I'm really impressed foundation DB folks. The talk "Testing Distributed Systems w/ Deterministic Simulation" by Will Wilson blew my mind. TL;DR they spent the majority of their initial dev effort into making a simulation of the database, then when they were happy with that plugged in real storage, time and networks at the end. Well worth a watch for anyone interested in distributed systems & reliability. https://www.youtube.com/watch?v=4fFDFbi3toc&pp=ygUgZGV0ZXJtaW5pc3RpYyBzaW11bGF0aW9uIHRlc3Rpbmc%3D https://www.youtube.com/watch?v=4fFDFbi3toc&pp=ygUgZGV0ZXJta...
- ergl 3y ago"f failures with f+1 replicas" is the standard for all non-byzantine fault tolerant systems out there. You will find it in Paxos, Raft, Viewstamped Replication, etc. It makes sense if you think about it: these systems follow a leader/replica model, and naturally you only need one leader to make progress
- robertlagrant 3y agoReplica is ambiguous here: is it 1 leader and n replicas? Or is it just n replicas, one of which is assigned "leader"? I thought "these systems follow a leader/replica model" would be the former, but "f failures with f+1 replicas" the latter.
- thwarted 3y agoIt's a cluster size of n replicas, with one of the n being the (current) leader. f failures with f+1 replicas is a cluster size of n replicas can sustain n-1 failures. n=f+1 or f=n-1. You wanna be able to sustain f failures, you need a cluster size (n) of f+1. When there is a failure, a non-failing node becomes the leader (or there's no leader change if the current leader isn't the one that failed). A cluster size of 1 has 1 leader, and can sustain 0 failures.
- angio 3y agoHow do people deploy FDB to the cloud? Is it possible to deploy it without EBS to take advantage of cheaper VM temporary storage?
- manishsharan 3y agoYes I would think so. FDB is distributed by default and the cluster is very easy to setup. As long as you have sufficient number of VMs in a cluster, the loss of a single vm or disk will not matter as you can spin up a new VM to join the cluster. On AWS, you can set up members of the cluster in different availability zones in the same region .. so the outage of on zone will not impact your database. I am running this set up in my dev (personal) environment on AWS.
- dialogbox 3y agoI won’t do that for production. Regional failure is not impossible although it is rare. You will lose all of your data.
- alberth 3y agoQueues Is it correct to assume FDB is the perfect framework for creating a queue?
- ramchip 3y agoWith FDB latency shoots up when a bunch of writers compete to update the same entry, because writers can have to retry many times before the write finally goes through, with a network round-trip each time. Personally I found this much harder to work with than a postgres queue with SKIP LOCKED for example. There's this however: "QuiCK: A Queuing System in CloudKit": https://www.foundationdb.org/files/QuiCK.pdf https://www.foundationdb.org/files/QuiCK.pdf. I suspect it really depends on what you expect from a queue, e.g. if you need strict FIFO or priorities, and how much effort you're willing to invest.
- mike_hearn 3y agoAn obvious question you face when deploying something like FDB is how to write your app on top of it. With FDB it's like RocksDB. You get a transactional key/value store, but that's a very low level interface for apps to work with. FDB provides "layers", such as the Record layer. It helps map data to keys and values. But a more sophisticated solution that I sometimes wish would take off is this library: https://permazen.io/ https://permazen.io/ It's a small(ish) open source project that implements an ORM-like API but significantly cleaned up, and it can run on any K/V backend. There's an FDB plugin for it, so you can connect your app directly to an FDB cluster using it. And with that you get built-in indexing, derived data, triggers, you can do queries using the Java collections API, there's a CLI, there's an API for making GUIs and everything else you might need for a business CRUD app. It's got a paper of its own and is quite impressive. There are a few big gaps vs an RDBMS though: 1. There's no query planner. You write your own plans by using functional maps/filters/folds etc in regular Java (or Kotlin or Python or any other language that can run on the JVM). 2. It's weak on analytics, because there's no access control and the ad-hoc query language is less convenient than SQL. 3. There's no network protocol other than FDB itself, which assumes low latency networks. So if there's a big distance between the user generating the queries and the servers, you have a problem and will need to introduce an app specific protocol (or move the code).
- dang 3y agohttps://permazen.io/ https://permazen.io/ hasn't appeared on HN before*. If you'd be willing to post it and then email hn@ycombinator.com a heads-up, we'll put the submission in the second-chance pool (https://news.ycombinator.com/pool https://news.ycombinator.com/pool, explained at https://news.ycombinator.com/item?id=26998308 https://news.ycombinator.com/item?id=26998308), so it will get a random placement on HN's front page. * and the only previous related submission appears to be https://news.ycombinator.com/item?id=21646037 https://news.ycombinator.com/item?id=21646037.
- mcsoft 3y agoWe have seriously looked at FoundationDB to replace our SQL-based storage for distributed writes. We decided not to proceed unless we are about to overgrow the existing deploy, a standard leader-follower setup on the off-the-shelf hardware. The limiting factor for the latter would be a number of NMVMe drives we could put into a single machine. It gives us couple dozen Tb of structured data (we don't store blobs in the database) before we have to worry. fdb is best when your workload is pretty well-defined and will stay such for a decade or so. It is not usually the case for new products which evolve fast. Two most famous installations of fdb are iTunes and Snowflake metadata. When you rewrite petabyte-size database in fdb, you transform continuous SRE/devops opex costs into developers capex investment. It comes with reduced risks for occasional data loss. For me it's mostly a financial decision, not really a technical one.
- Jgrubb 3y ago> transform continuous SRE/devops opex costs into developers capex investment Would you mind expanding/educating me on this point? When I think of capex I think of “purchasing a thing that’s depreciated over a time window”. If you’d said “transform SRE/COGS costs into developer/R&D/opex costs” I would’ve understood, but eventually the thing leaves development and goes back into COGS.
- mcsoft 3y agoI assume a couple of things here: 1) that SRE costs would be lower with fdb at scale due to its handling outages, i.e. auto-resharding; and 2) that a migration project from *sql to fdb will be finite (hence an investment I hastily called capex). Would love to hear from anyone with experience in fdb whether these assumptions hold.
- foobiekr 3y agoBasically the SREs don't have anything to do with fdb for the most part. You add a node, quiesce a node, delete a node. Otherwise it's self-balancing and trouble-free from an SRE pov. See my other message for the developer issues, though. IMHO fdb as it is today is too hard for most developers if their use case is anything beyond redis simple keys.
- sproketboy 3y ago[dead]
- zinodaur 3y agoIf theres just one Sequencer, and every ReadVersion request to the Proxy eventually hits the Sequencer 1-1, how does the Sequencer not get crushed? Or is a scaling limit just "the number of ReadVersion requests a Sequencer machine can handle per second", which admittedly is a cheap request to respond to
- lowbloodsugar 3y agoYeah that seems like an untenable design choice. Was quite interested until I read that. Max TPS? and MTTR when sequence inevitably shits itself?
- richieartoul 3y agoReplied above
- foobiekr 3y agoYou can trivially scale fdb to tens of millions of tx/sec for write-heavy workloads without a hardcore cluster for transactions of reasonable complexity (though with careful design on my part and the part of others for collisions to be unlikely). MTTR on failure is seconds. Really, there's no system I've used that is as robust and performant as fdb and I include s3 in that list - s3, for example, _routinely_ has operations with orders of magnitude latency variance and huge, correlated spikes.
- richieartoul 3y agoRequests to the sequencer are batched heavily. If the sequencer fails, the cluster goes through a recovery and will be unavailable for 2-3 seconds and then recover.
- zinodaur 3y agoGood point about the batching! Any idea what kind of ReadVersion qps throughput you can get this way? And yeah, 2-3s unavailability seems fine.
- 3y ago
- pi-r-p 3y agoIn my company, we tested FDB for two years, then we wrote a new backend for Warp 10 timeseries database... Performances are really impressive, we dropped HBase backend when we released Warp 10 3.0. Note we can isolate customers easily on the same FDB cluster (tenants are not explained anywhere on internet, it is a quite recent FDB feature). more info: https://blog.senx.io/introducing-warp-10-3-0/ https://blog.senx.io/introducing-warp-10-3-0/
- foobiekr 3y agoWe have run foundationdb in production for roughly 10 years. It is solid, mostly trouble free (with one very important exception: you must NEVER allow any node on the cluster to exceed 90% full), robust and insanely fast (10M+ tx/sec). It is convenient, has a nice programming model, and the client includes the ability to inject random failures. That said, I think most coders just can't deal with it. For reasons I won't go into, I came to fdb already fully aware of the compromises that software transactional memories have, and fdb roughly matches the semantics of those: retry on failure, a maximum transaction size, a maximum transaction time, and so on. For those who haven't used it, start here: https://apple.github.io/foundationdb/developer-guide.html https://apple.github.io/foundationdb/developer-guide.html ; especially the section on transactions. These constraints _very_ inconvenient for many kinds of applications so, ok, you'd like a wrapper library that handles them gracefully and hides the details (for example count of range). This seems like it should be easy to do - after all, the expectation is that _application developers_ do it directly - but it isn't actually so in practice and introduces a layering violation into the data modeling if you have any part of your application doing direct key access. I recommend people try it. It can surely be done, but that layer is now as critical as the DB itself, and that has interesting risks. At heart, the problem is, the limits are low enough that normal applications can and do run into them, and they are annoying. It would be really nice if the FDB team would build this next layer themselves with the same degree of testing but they themselves have not, and I think it's pretty clear that it turns out a small-transaction KV store is not enough to build complex layers in actuality. Emphasis on the tested part - it's all well and good for fdb to be rock solid, but what needs to be there is that the actual interfact used by 90% of applications is rock solid, and if you exceed basic small-size keys or time, that isn't really true.
- Dave_Rosenthal 3y agoI think that’s a good and really fair summary. - If you’re a developer wanting to build an application, you should really use a well designed layer between yourself and FDB. A few are out there. - If you’re a dev thinking you want to build a database from scratch you probably should just use FDB as the storage engine and work on the other parts. To start, at very least! (One last thing that I think is a bit overlooked with FDB is how easy it is to build basically any data structure you can in memory in FDB instead. Not that it solves the transaction timeout stuff, etc. but if you want to build skip list, or a quad tree, or a vector store, or whatever else, you can quite easily use FDB to build a durable, distributed, transactional, multi-user version of the same. You don’t have to stick to just boring tables.)
- tiffanyh 3y agoFDB uses SQLite to store data, but FDB doesn’t expose SQL to the end user. https://apple.github.io/foundationdb/architecture.html https://apple.github.io/foundationdb/architecture.html
- richieartoul 3y agoNewer versions are moving towards a custom btree storage engine called Redwood
- Dave_Rosenthal 3y agoThat’s true, but FDB doesn’t use very much of SQLite, just a modified version of SQLite’s internal b-tree.
- AtlasBarfed 3y agoThe Sequencer: - does not have a persistent/disk-backed state - It is a singleton process - it and only it does order, no logs do ordering ... if the singleton sequencer crashes, I do not see on this high level description how the system recovers, if the sequencer is the only one that knows write order but has no persistent write "log". What am I missing? This... does not appear to be something you run outside of a dedicated datacenter, AWS with its awful networking and slow/silently throttling storage would probably muck this thing up under any substantive scale?
- richieartoul 3y agoIt runs fine in AWS, Snowflake and many others run it there. The most recent FoundationDB paper goes into a lot more detail on their recovery protocol, it’s a lot more nuanced than you think, but it works extremely well
- Dave_Rosenthal 3y agoWhat you are missing is that the "tlogs" (transaction logs) actually hold the durable, fault tolerant write log. The sequencer is just a big fast in-memory data structure that checks if the many transactions coming into the system pass isolation checks (the I in ACID). That is, it accepts transaction so long as the keys that the transaction read haven't been modified in the mean time. The reason it can fail without a correctness issue is that it can just reject all transactions in flight for the clients to retry. This is something the clients need to be prepared to do anyway because of optimistic concurrency. It can run fine on AWS. Upon a failure, the sequencer role is very fast to re-elect onto another machine in the cluster because there is no persistent state at all.
- chalcolithic 3y agoare there dumber alternatives? sort of like ndb cluster(its kv interface, to be precise) but fully disk based, so that transaction limits are practically unreachable?
- louisefoster 3y ago[dead]