Google and Red Hat announce cloud-based scalable file servers
googlecloudplatform.blogspot.com
googlecloudplatform.blogspot.com
Ex-GlusterFS person here (used to work at Red Hat on the project side, leaving mid last year).
"Small file access", and "lots of files in a directory" have been a pain point with GlusterFS for ages. The 3.7.0 release had some important improvements in it, specificially designed to help fix that:
https://www.gluster.org/community/documentation/index.php/Fe...
The latest Gluster release is 3.7.8 (the same series as 3.7.0), and is worth looking at if you're needing a good distributed file system. If you have something like 1Mill files in a single directory though... hrmmm... NFS or other technologies might still be a better idea. ;)
The basic answer to your question is that networks are slow compared to local storage. In order to get decent performance, you must either avoid network round trips or amortize their cost over many operations. We - not just GlusterFS but distributed filesystems in general - can do this pretty well for some operations. We can batch, buffer, cache, etc. This works great for plain old reads and writes to large files (for example). It doesn't work so well for operations that have to touch many small files. For those, in order to ensure the required level of consistency/currency, we have two choices.
(1) Send a request per file to get current metadata.
(2) Cache metadata, and participate in some sort of consistency protocol to make sure we don't serve stale cached information.
Both approaches have workloads where they perform better and workloads where they perform worse. In addition, the second approach adds a lot of complexity, especially in a system where failures are common and a loosely coordinated set of servers must respond (systems with a single master server have an easier time here but are less resilient). The inherent difficulty of this approach is why e.g. CephFS has taken so long to mature.
GlusterFS took the first route instead. It does mean that "ZOT" (Zillions Of Tiny files) workloads will perform poorly. I won't deny that. On the other hand, it's easier to test or prove correct, and the time not spent on solving the hard version of that problem - often for little eventual benefit - can instead be spent on other kinds of improvements. Some people are happy with that. Some are not. Some spread FUD. Some try to implement the practical equivalent of a distributed filesystem on top of some alternative (e.g. object stores) with their own even more serious limitations, and experience even more pain as a result. Some are initially unhappy with these tradeoffs, but work with us and learn to work more effectively with these limitations to enjoy other benefits. That's life in the big city.
But I don't know shit about filesystem design, soooooo...
Unfortunately, I don't remember most of the details. It's fallen out of my headspace in the months since I've left. ;)
It was the whole shebang: Kernel panics, inconsistent views, data loss, very slow performance, split-brain problems all the time. Our set up IIRC was very simple: two bricks in a replicated volume. It worked so poorly that we had to take it out of production. Some of our experience can be explained by GlusterFS performing poorly under network partitions, but nothing could justify kernel panics. It blew my mind that Redhat acquired that company and product.
Edit: I hope there's been a big improvement to the reliability and performance of GlusterFS. Can anyone with more recent experience running it in production comment?
In a 3-node cluster, any system with a decent consensus algorithm (to be clear, I'm not sure if GlusterFS has one) would know that during a partition the cluster can only continue to operate if at least 2 nodes can communicate with each other to elect a new master.
One has to wonder why the software even lets you run it with an unsafe number of nodes though... (Lots of distributed software is guilty of this...)
Still, if I can't make it work well on a reasonably small scale, why would I expect it to work any better when I scale it up?
This is not a given, the cluster can and should refuse to operate if it can not get a majority vote for the master, if that is required to prevent data corruption. With two nodes that means both nodes must be active and reachable, with three nodes one node can fail.
As far as performing poorly under network partitions, I'd love to hear more. That is our responsibility, and sounds like something we can/should fix.
If I had to make a comparison, I'd say GlusterFS reminds me a lot of MongoDB in the beginning. It wins a lot of kudos at the outset based on ease of setup, management and CLI UI, plus it has a good "story" on ability to scale up that gradually begins to fray when pushed. Hopefully there have been big improvements.
Do you mind sharing which GlusterFS version you're on, which kernel, and roughly what your load profile is? (eg. lots of small files, read-heavy, write-heavy, etc.) Also, are you using Red Hat Gluster Storage or the open source version?
Ubuntu 14.04.3 LTS 3.13.0-24-generic g
read heavy with a mix of 100MB or 15GB files depending on datasets
Unfortunately, I hit a roadblock in relation to enumeration of huge directories: Even with just 5K files in a directory, performance started to drop really badly to the point where enumerating a directory containing 10K files would take longer than 5 minutes.
Yes. You're not supposed to store many files in a directory, but this was about giving third parties FTP upload access for product pictures and I can't possibly ask them to follow any schema for file and folder naming. These people want a directory to put stuff to with their GUI FTP client and they want their client to be able to not upload files if the target already exists. So having all files in one directory was a huge improvement UX-wise.
So in the end, I had to move to nfs on top of drbd to provide shared backend storage. Enumerating 20K files over NFS still isn't fast but completes within 2 seconds instead of more than 5 minutes.
Of course, now that we're talking about GlusterFS, I wonder whether this has been fixed since?
It maybe worth a second look - however a directory with a large file count still might be a problem
I don't. I am using stock vsftp with a PAM module that allows authentication against our web application.
Gluster is available not announcing that this is a new technology.
I think, though I'm not sure, that its an official, active support relationship between Google and Red Hat for GlusterFS on GCE, related to Red Hat's vendor certification programs.
If so, this seems like the kind of thing that has the most impact on people with support contracts with Red Hat, but should also have peripheral impacts on the stability, quality of support, etc., of the product on the GCE platform generally.
Overall, I feel AWS S3 is a better (or at least simpler) approach. Just acknowledge that files are not locally stored and use them as is. AWS is experimenting EFS as well, which we found not as desirable as well.
Edit: I am not saying that you cannot make GlusterFS or EFS perform great. My appoint it that it's hard to do so, and might not worth the effort to develop such a system given that S3 can serve most needs of distributed file storage.
My solution was to mount an EBS on the "worker box," along with an NFS server. Each "ingress box" runs an NFS client that connects to the server via its internal VPC IP address, and mounts the NFS volume to a local directory. It works wonderfully. In three months of running this setup, I've had no downtime or issues, not even minor ones. Granted I don't need any kind of extreme I/O performance, so I haven't measured it, but this system took less than an hour to setup and fit my needs perfectly.
We've been using it in production for a few years now and having a single namespace that can basically grow ad infinitum has been pretty neat.
If you want a trouble free Gluster experience stay away from MANY small files and replicated volumes.
Also, there's a sentence on the end of the GlusterFS section there saying:
"If you want to deploy a Red Hat Gluster Storage cluster on Compute Engine, see this white paper for instructions on how to provision a multi-node cluster that includes cross-zone and cross-region replication:"
Apart from the typo (pedant alert!), the URL on "white paper" goes to a non-public document only Red Hat subscribers have access to. That should probably be that fixed, so non-Red-Hat-subscribers can read the doc and know what they'll need to do up front.
If people need to subscribe to RH in order to get that info... just to know what they need to do... that's probably going to hinder adoption. Potentially by a lot. ;)
Is anyone here running a GlusterFS setup with high read/write volume on small files successfully? If so, what's your secret?
http://www.highlyscalablesystems.com/3202/colossus-successor...
So, here is the link to the star of the show: https://www.redhat.com/en/technologies/storage/gluster
We designed it to yield top performance both for (parallel) file system workloads and block storage workloads. So you can run VMs and databases on it with a performance better than any other partition-tolerant software storage system. The goal is to provide customers with a scalable automated general storage platform for all workloads, a la Google, but for real world applications.
The issue with Lustre is that it usually requires kernel patches if you are not running the aforementioned distributions that are designed for HPC. Specifically the Lustre OSD used a modified version of ext2 among other issues.
That could be getting better now with the new ZFS based OSD though I wouldn't put money on it being "easy" to install.
Lustre also doesn't provide any means of replication. If you want to achieve HA with Lustre you need to make each OSD individually HA. This can be done with multi-pathed SAS arrays and a ton of scripting but it's still not exactly a walk in the park.
Hopefully one day we will see a real high performance distributed filesystem that also bundles replication, tiering and some semblance of POSIX compatiblity. I doubt Ceph is going to be it so we are probably still 5-10 years from a solution to the problem.
I'm not so sure the hope for "one filesystem to rule them all" will ever really work out, but Ceph is the best positioned.
For example, what if you have to run code that depends on a POSIX interface? That "just" becomes a massive rewrite project.
Disclaimer: I work on Compute Engine.
Really looking forward to EFS to take a big bite out of that commitment since it has reasonable pricing for projects with smaller storage requirements.
Just switching a portion (photos) of the project with the most files to S3 took on the order of 800 (mostly developer) hours.
At some point you have to decide what you can do with the hand you're dealt and wether you want to distribute those tens of thousands of dollars to your employees, pursuing new business or your own pocket; or reinvest it into an existing project despite having no mandate to do so. It's not an easy or cut and dry a decision to make (even today) for legacy projects IME.
If you're starting a new project. Yes. S3 is the most reliable service in the AWS stack (AFAIK) by a big margin. The pricing is very reasonable today unless you're in the petabyte range maybe, and you have no operations overhead. Agreed it's a no-brainer.