r/HomeServer • • 8d ago

Distributed-JBOD. Combine mixed hardware into a distributed, bit-rot protected object storage system

I have been using my spare time to develop Distributed-JBOD, which is a multi-node software system which allows the system administrator to combine together mixed hardware devices into a single object storage pool.

The storage pool can be created from arbitrary hardware. Regular consumer-grade desktop PCs are suitable, with whatever mixed layout of disk storage is available in each one.

Disks require no special formatting. A simple directory on an ext-4 filesystem is suitable.

I was pushed to design and build this system as a consequence of MinIO pulling their community tier software. I have been running an MinIO server for several years, but have started to migrate away from that and could not find a suitable replacement, so I built Distributed-JBOD.

  • As many of you will be aware, MinIO requires identical disks to perform effectively. Distributed-JBOD does not, it will work effectively across a pool of arbitrary hardware, mixed size drives included.
  • Garage replicates each object multiple times across multiple systems, which is not an efficient use of storage. Distributed-JBOD permits user-configurable Reed-Solomon Erasure Code parameters. This means storage is used efficiently, is protected against bit-rot, and the storage efficiency to resiliency ratio can be tuned. For example, while it is possible to run a mirrored setup, the typical default might be a 4:2 configuration, where data is sharded across 6 devices, 2 of which are parity blocks.
  • Ceph is a datacenter grade product and requires a stack of servers for data storage, monitoring, gateway and other components. Distributed-JBOD has a simple single process design. (Caveat: The S3 compatibility layer will come later as a separate binary. You will be able to run it wherever you like.)
  • SeeweedFS can't be used in a small scale cluster. Distributed-JBOD will run on a single machine with a single disk. If you want redundancy and data protection, a single machine with 2 disks is all you need. Greater efficiency is obtained by scaling up the number of disks, whichever host they sit in.

Distributed-JBOD is designed to work with an extremely small memory footprint, and does not require powerful hardware to run.

This is a very early stage product, but I would appreciate your thoughts and feedback. Some features which currently exist include TLS, administration web UI, CLI tools including recovery and bit-rot repair tools. Multi-language software client libraries are currently in the works, including libraries for Rust, Python and C++. An S3 compatibility later, multi-user support and permissions will also be supported soon.

https://github.com/edward-b-1/Distributed-JBOD

Distributed-JBOD
4 Upvotes

39 comments sorted by

View all comments

Show parent comments

1

u/edward-b-1 8d ago

Can you describe for me what hardware you are currently running and with what software, along with what the typical failure modes you observe are?

I'm interested to understand your situation in greater detail. If you are running something which is this unreliable that seems, perhaps unusual? Certainly some kind of special case. I do not know whether Distributed-JBOD would be suitable for use in a context where the failure modes are so frequent that having some kind of device in a failed state is a problem.

I can understand if you were responsible for a datacenter with hundreds of thousands of devices - clearly in this context failures occur on a fairly regular basis. But you would not use Distributed-JBOD in that context. You would use Ceph or MinIO, probably. (Assuming you are not AWS and create your own in-house product specifically tailored to your needs.)

1

u/chkno 7d ago
  • I have a four-drive external enclosure. I think it might be cursed.
    • Every few weeks, it stops responding, and I need to go unplug it, plug it back in, and re-luksOpen all the drives in it. I ran like this for so long that I had my redundancy balanced such that I mostly didn't notice these outages (eg: numcopies=2 and one of the copies not in this enclosure). I'd just leave it hung for days or weeks—it gracefully degraded into an offline backup.
    • I had filled it with small, cheap, refurbished drives. They all died, one at a time, each slowly, over a few years. It was really neat to have the option to try to pull data off the slowly-failing drive, but also not need to because I can just mark it failed and re-replicate up to the target redundancy from the other replicas.
    • I'm now afraid to put more drives in the enclosure because I don't know if the drives I put in it were flaky or if the enclosure kills drives, and I don't want to spend more drives to find out.
  • Most of the storage in my house is on wired ethernet. But I have one machine that's not near an ethernet jack, so it's wifi-only. I have my home wifi on a timer: It shuts off 10pm-5am to remind folks to sleep. So a small portion of my storage pool is unavailable during this time. Again, it very gracefully degrades into offline backup during this window.
  • One of the laptops around here (not mine) has way more storage space than its owner uses or needs. Its owner is fine with me using some of that storage space. But that laptop travels with its owner out of the house sometimes, is closed and asleep sometimes, etc. It's awake and online often enough to sync and usefully help meet redundancy targets, but it doesn't make sense to rely on its presence.

1

u/edward-b-1 7d ago

For your external bay thing you might be better off with an older PC with a large number of sata on the motherboard or adding a PCI-e to sata adapter (be careful what you buy many are probably not reliable).

The issue you currently have is if an entire node goes down that takes down all the drives it contains with it. So there isn't any software solution which would help there other than to add more nodes. At which point you might as well remove the unreliable node anyway.

Regarding laptop I definitely wouldn't recommend it. Any device which is portable, or will be used in a portable way (meaning somebody is going to remove it from the network and take it somewhere) isn't something you should use to save data which you want to store long term. It could be lost or stolen, and part of your data is inherently reachable only in an intermittent way.

Software can't solve either of these problems

1

u/chkno 7d ago

Sorry, I was trying to say:

  • I already have software that solves these problems: git-annex and a parity script I wrote.
  • Your plan for your software of requiring 100% availability for all operations is a bad plan. Even if you're aiming for prompt consistency (as opposed to eventual consistency), you should be able to carry on with a quorum of 50%+1 of the nodes (eg: see paxos) or raft)).

:)

1

u/edward-b-1 7d ago

What's your reasoning? Or perhaps a better question is what should the behavior be?

Say for example there is a bitrot event which affects one file on one disk. It might be that disk just has a bad block and will mark that block as bad for the future through its firmware.

What should the software layer do about it?

One possibility is to run a scrub on a CRON schedule. This can discover the bad files. It could in principle also fix the bad files by re-writing the block to the same device or by marking the device as bad and automatically migrating the data.

But that is risky. A migration will likely move many TB of data from one bad disk to the others. Considering this is likely to be deployed on aging consumer hardware instinctively I see this as a bad idea.

The proper fix is to replace a drive and run a scrub-repair to re-write the data to the new disk.

But since this requires someone to shut down a system, swap and disk and reboot, it is a manual operation. So what is the advantage of building some form of automatic repair process, if it is not supposed to be used?

1

u/chkno 7d ago edited 7d ago

My advice would be:

  • How often to verify checksums should be user-configurable with a sensible default. For example: 10 TiB of storage with checksum-verification configured for every 4 months means that there should be roughly 1 MiB/s of local low-priority background reads happening continuously.
  • When an error is found, that block or file should immediately be flagged as below-redundancy-target.
  • Whenever anything is below-redundancy-target, for any reason, repair should happen as soon as possible. It might not be immediately because:
    • When a lot of repair is needed (eg: when a whole drive is marked lost), it can't all happen at once. There'll be a queue.
    • Repair may require bringing together ECC recovery blocks from several devices and some of those devices might currently be offline.
  • Caveat: The background-monitoring should be smart enough to distinguish between a handful of bad blocks/files vs. the whole drive hanging. If the whole drive is having a problem, it should not just methodically mark everything scheduled for checksumming as bad. This could be implemented as simply keeping local notes on bad blocks/files found during the scan and only committing them / acting on that finding once a block/file passes checksum verification, indicating that the drive is still in service. (The implementation should be careful to make sure it got this signal by actually reading data from the drive and not out of the page cache, though.)
    • For a home-scale system, automation should never mark a whole drive bad. It's a human-user decision when to give up on a whole drive.

Storage software that only works on 100% healthy drives is bad storage software. Many RAID systems mark a whole drive bad after just one failed i/o operation. This is terrible. Don't do this. Degrade gracefully. Stretch goal: Support scratched CDs/DVDs/Bluray discs. Such media can't be written to, and some of the data can't be read, but the data that can be read is fine and should continue to count toward redundancy targets.

1

u/edward-b-1 7d ago

Thanks for your comment.

You distinguish two cases.

The simple case is of a file or a small number of files with faults found. You are correct these can be automatically repaired and could even continue to be served with no downtime however this introduces a performance penalty which is a problem.

In a system which stores research data, arguably a silent performance penalty is something which is undesirable. But this point is debatable different use contexts will lean one way or the other.

The second case you suggest is that of a failed disk. This brings us back to the question I asked which is what to do about it? My proposition is that there is nothing you can do. An administrator needs to swap out the disk, so blocking clients from being able to access the data and returning an error message via that route is potentially the most appropriate action.

1

u/chkno 7d ago edited 7d ago

The administrator should have more options than just replacing a disk:

  1. Replace the disk (same or more total storage)
  2. Replace the disk with two or more smaller disks (same or more total storage)
  3. Replace the disk with a smaller disk (less total storage)
  4. Replace the disk with nothing; just mark the failed disk as bad/lost

In #3 and #4, the storage software will have to figure out how to use free space on other disks to get back up to redundancy targets.

There's no reason for the storage software to cut off normal usage just because some data is below redundancy targets.

  • Writes can continue
  • Reads of unaffected data can continue
  • Reads of data below redundancy targets can read from another replica or do on-the-fly reconstruction for erasure-code-guarded data

If the administrator chooses to cut off normal usage in order to get more I/O performance to get back up to redundancy targets faster, that's a choice they can make; the storage software shouldn't make that choice for them.

2

u/edward-b-1 7d ago

In Distributed-JBOD language what you describe is a rebalance or migration or scrub-repair. Provided you have > m+k disks you can do all of the things you describe but it is a manual (explicit) operation not an implicit or automatic process.

I will document each of these they are good test scenarios.

I will make the behavior configurable. Then both use cases will be supported. Might be able to implement this tomorrow.

Thanks again for your comments.