Why do you think there's no open-source equivalent?
Why do you think there's no open-source equivalent?
We've been told for years that file system primitives are bad, unscalable and nasty. part of this is because NFS can be a dick, but a lot is because large scale file systems generally are fast, scalable, coherent, secure and reliable (pick two).
A lot of things can be managed with S3, but it has a lot of flaws.
Unless you've worked at a place where you've had access to something like GPFS(or whatever IBM calls it now), where you can put hooks into everything, and programme how you want to mirror/cache/copy namespaces, it all seems impossible.
Take FB for instance. They've gone balls deep with an S3 clone, which is both slow and difficult to use. Every operation takes about 3 seconds. (pulling a 10k file? pulling a 100meg file? list operation? mkdir? all 3 seconds.) Its an interface on top of the underlying blob store, so its looks like it stores it's metadata in a SQL Db (ie like lustre with metadata node and object storage nodes)
That sucks for machine learning or any kind of "normal" file operations this means there is lots of heavy engineering to get something like performance that would be solved by an ephemeral clustered file system.
The other big thing to note is that clustered filesystems are a massive massive pain in the arse to manage, especially if you choose the wrong one. (looks knowingly in the direction of ceph and gluster) you need to know what tradeoffs you are choosing before you install. Lustre is/was effectively a raid-0 over the network, Ceph is just plain inefficient, gluster is just silly, GPFS is held back by IBM being IBM.
for most people clustered filesystems are wrong choice.
Do I understand correctly that this would be one possible interpretation of this sentence: If you've only used UNIX, and not something IBM made decades ago, then it seems impossible.
That's hilarious.
The general gist is that unless you've worked in a place that has setup a "proper" network, that is roaming home directories, mapped posix storage, kerberised user account, and some sort of multi-machine graph based job dispatch system (ie airflow, grid engine, large scale k8s with a batch plugin, aws's batch pixar's tractor), then you will have not seen a need for a large clustered filesystem with a posix-like interface.
moreover, given the "filesystems are hard and bad, use object storage instead" noise that we've had since we joined the cloud era, why would you _try_ to use a clustered filesystem? Especially as most of the "free" ones are shite.
1. The abstraction is not what most developers want. Colossus gives you append-only files, halfway between a file store and an object store. This is not an abstraction that most developers are comfortable with. It also has a relatively thick client library, so if you were planning to have a "microservice" it is a lot to import. Both of these factors help to simplify the reliability story.
2. Open source users are generally used to picking the abstraction they want, then finding the thing that does it. In contrast, Colossus/GFS use is mandated at G.
3. Colossus's bugs were worked out due to scale, which means wide adoption. Ceph and Lustre, in contrast, tend to have relatively small installations (measured in TiB or single-digit PiB) and their maintainers have seen many fewer byzantine failures by consequence.
I have played with a SaaS offering of a similar service, but it's a huge project and I don't want Google to sue me. OSS would likely be completely off the table except as an "open core" type of concept - there is just too much engineering effort needed to get to the first version.
Most of the files I use are append-only (or in fact write-once).
The only files that are random-access are database files (but I prefer using the filesystem as an approximation to a database; not sure to what extent Colossus can serve that usecase; e.g. does it provide transactions?)
Do you have other examples of files that are not append-only or write-once?
Another example is documents (think .doc files), although these can also be log-structured or checkpointed (I think the .docx is almost 100% log-structured).
- SaaS solution, not just open source (improvements to s3, for example)
- A global hierarchy of files rather than an assortment of different buckets like s3
- Per-subdirectory permissions and chargeback so anyone can create project-specific subdirectories and have their organization billed for storage costs
- Filesystem drivers allowing every production and development computer and coding environment to access remote files in the same way and with the same ease as remote files. The difference between accessing remote files and local files in Google code is literally just the file prefix. In all other ways, remote files are _better_ than local files (they have all the same features, plus more) so there's very little reason to directly use local disk most of the time.
Most companies don't need to process hundreds of terabytes per second to and from disk and so they can centralize storage behind a few network interfaces. Centralized storage is conceptually simpler and more straightforward to manage.
In short, very few folks have the architecture to run something like Colossus, or the need for it, so it doesn't attract open-source replication. Gluster, Ceph, etc. are the closest.
Two reasons:
1. Colossus is extremely dependent on other Google Tech to work and be cool. Implementing 'just colossus' without the rest of the tech stack would require developing quite a diverse set of other tools first.
2. Part of what makes it cool is that everybody is mandated to use it, and hardware, kernel and so forth improvements are made to make it better.
At Quobyte we're building a scalable fault-tolerant file system (POSIX, so it has a different architecture and constraints), but it's closed source.
I've tried a couple times at different companies to get developers to stop assuming POSIX and switch to object storage as their persistent byte abstraction layer.
Trying to scale POSIX is a dead-end. Every time I try and see people build it, it's a big pile of suck. For example, CephFS. Yea, it works, but it still sucks compared to object storage. There's just no magic in the universe that will get you around the CAP issues with trying to make POSIX work over a network.
If you go full object storage semantics (no in-place updates, to appends, no rename, ..) you're pushing the problem to the application layer. Experience at Google was that that's not a good place to have it, because application developers are usually not good at solving distributed systems problems. This lesson learned is the part of the Bigtable (no transactions) to Megastore (application layer transactions on top of Bigtable) to Spanner (transactions built-in) journey.
So yes, of course there are inherent hard trade-offs, but dropping strong semantics from the storage layer by using an object store instead of a file system is usually not a good idea.
Also, there don't seem to be public research papers about it.
The question of atomic clocks aside (whose answer I don't know), Ex-Googlers have created open-source alternatives to Google-internal software in a lot of other cases, too, so I don't think the lack of public research papers could be the reason.
The atomic clock thing helps improve performance of Spanner, by tightening the SLO of node-to-node and continent-to-continent clock drift.