— ceph, bluestore, storage, distributed-systems — 12 min read
Hey folks, as a continuation of my goal to understand BlueStore more deeply, I am writing my second blog post, revolving around the question: "Why is BlueStore the way it is?" When reading the BlueStore code and resources, the following questions kept bugging me:
I personally understand concepts better when I get some historical context that led to the present. The quote "Those who cannot remember the past are condemned to repeat it" forms one of the bases of my learning structure. This is what led me on the quest for the origins of BlueStore. Fortunately, I stumbled across the paper titled "File Systems Unfit as Distributed Storage Backends: Lessons from 10 Years of Ceph Evolution" (Aghayev et al., SOSP '19).
I encourage everyone to read the paper. This post is hugely motivated by the paper and thus contains some excerpts from it.
In this series we will explore the evolution of the storage backends that were used in Ceph, how the problems with each of them motivated the next version, and how all of the learnings combined paved the way for BlueStore. This turned out to be a long one, so I have split it into three parts:
A storage backend is defined as the software module directly managing the storage device attached to a physical machine, i.e., it's the software on each machine that actually puts the bytes on the disk. These backends are entirely local. In a distributed system, where there are many physical storage devices, it is the role of the system to aggregate the storage backends on the many physical devices and present them as a single unified data store. Note that the distributed layer's guarantees are built on top of the backend's guarantees.
Side note: when the paper says "Distributed Storage Backends", read it as (Distributed Storage) Backends, i.e., the backends that are used in the distributed storage system.
1┌──────────── Ceph / RADOS (distributed system) ────────────┐2 │ CRUSH, PGs, replication, recovery, cluster maps │3 └───────┬─────────────────┬─────────────────┬───────────────┘4 │ │ │5 OSD.0 on node A OSD.1 on node B OSD.2 on node C6 │ │ │7 BlueStore #0 BlueStore #1 BlueStore #2 ← "distributed storage backends"8 │ │ │ (three independent, local instances)9 /dev/sdb /dev/sdb /dev/sdbIn the terminology of Ceph, RADOS can only promise strong consistency if each OSD can apply a multi-op transaction atomically.
Why? Because a single client write is a transaction that updates several things at once:
When a crash happens, RADOS checks which replicas are up to date by comparing the logs. This makes sure that the client doesn't receive stale data. This is only possible if each OSD's log tells the truth about its data. If a crash could leave the log entry for version 42 on disk while the data is still at version 41, the OSD would claim that it has the current data while still holding stale bytes, and RADOS would believe it, serving the client old data.
Hence the backend must apply each transaction as all-or-nothing on its own disk. Keeping replicas in agreement is RADOS's job. Each backend must just make sure that what it reports about itself is true.
It is my understanding that a good backend has the following properties:
Each of these properties is hard to get on top of a file system, and this maps nicely to the three problems we will look at later in this post:
| A good backend needs | Which is hard on a file system because of |
|---|---|
| Strong consistency | Transactions: there are no multi-op atomic transactions |
| Efficiency | Transactions (costly to fake) and Metadata (slow with millions of objects) |
| Longevity | New hardware: we have to wait for the file system to support it |
Historically, storage began with people/companies owning the disk. As one might imagine, it caused a lot of issues. Later, when local file systems were developed, storage adopted them, as they were mature, familiar and often good enough for the large-file workloads of the 2000s. With changing workloads (small objects and transactions) and advancements in hardware (NVMe, zones), we are outgrowing what a general-purpose file system can offer cheaply and are returning to owning the disks.
Ceph's own path is the whole story in miniature:
We'll explore each of them in detail in Part 2.
File systems became the de facto standard for the following reasons:
ls and find can be used to explore the
disk contents.For about 10 years, after a brief start with EBOFS, Ceph built its storage backend on local file systems (FileStore on Btrfs → FileStore on XFS → NewStore).
There are three main problems with using file systems as the storage backend for distributed systems:
We will explore each of these issues in this post, and in Part 2 we will map which backend had which issues and how they motivated the design of BlueStore.
As we saw in the earlier section, a distributed system's consistency guarantees are built on top of the guarantees of its storage backends. This means each backend must be able to apply multi-op transactions atomically, so that after a crash the local state is the same as the one it reports to the cluster. Hence, atomicity (all-or-nothing) during a transaction is a hard requirement.
POSIX file systems have no transaction APIs. This means there is no way to
express something like "apply these three operations atomically". They give us
individual atomic operations (rename, open, etc.) but no multi-op
transactions, and durability comes only from fsync. File systems do use
transactions internally to keep their own metadata consistent, but generally
don't expose them to applications.
For example, in order to atomically replace a file, we need to do something like this:
1// 1. Write the new version to a temporary file2int fd = open("config.tmp", O_WRONLY | O_CREAT | O_TRUNC, 0644);3write(fd, new_contents, len);4
5// 2. Make the new contents durable6fsync(fd);7close(fd);8
9// 3. Atomically swap it into place10rename("config.tmp", "config");11
12// 4. Make the swap itself durable13int dirfd = open(".", O_RDONLY);14fsync(dirfd);15close(dirfd);Notice that after each step we call fsync to make it durable. This is fine if
there is only one thing to switch, i.e., the directory entry that maps the name
config to an inode. A storage backend like Ceph's needs a single write to
update the object's data, its metadata and a log entry, all together.
One cannot do the following:
1op1: write object data (v42 bytes)2fsync3op2: set xattr version = 424fsync5op3: append PG log entry 426fsyncIf a crash happens at any of these steps, the intermediate data is wrong, and it's hard to roll back.
Hence a backend that is built on a file system has to construct transactions on top of it. There are three ways to do it, and Ceph tried all of them:
The next sections talk about why these options have significant performance/complexity overhead.
Many file systems have an in-kernel transaction framework. This is a mechanism the file system uses so that its own multi-step updates are crash-safe. For example, when you create a file, the file system updates the free-space bitmap, the inode and the directory entry, and wraps them in an internal transaction so that a crash can't leave them half-done.
A sample in-kernel transaction (from ext4's journal, jbd2) looks like:
1handle = jbd2_journal_start(journal, nblocks); // reserve journal space2 // modify bitmap block in memory3 // modify inode block in memory4 // modify directory block in memory5jbd2_journal_stop(handle);The idea is to expose this machinery to user space so that applications like Ceph can get atomic multi-op transactions for free.
Why these internal transactions aren't enough
These transactions were built for a different job, i.e., to keep the file system's own structure consistent. Hence they have the following properties:
The property that bites Ceph is "no rollback". Let's say we expose these internal transactions for applications to use. We could then do something like:
1TRANS_START2 op1 ✅ applied3 ← process dies4 op2 never issued5 op3 never issued6TRANS_END never issuedSince this is an in-kernel transaction started from user space, the kernel now holds an open transaction whose owner is gone. It cannot undo op1, because it has no rollback mechanism. Its only option is to close the transaction as-is and commit op1 alone.
This leads to a state where the object on disk has v42 data but a v41 version and a v41 log entry, which is an inconsistent state.
But why can't the kernel just undo op1? This is because file system journals
like jbd2 are redo-only, i.e., they only record the new contents of each block
and never the old ones. So there is nothing to undo with. Rollback was never
needed in the first place, since kernel operations validate everything up front
and don't fail halfway. If something truly unexpected happens in the middle of a
transaction, the file system does not roll back; it aborts the journal and goes
read-only (e.g., ext4's errors=remount-ro, a Btrfs transaction abort).
Note that this only bites when the process dies and the kernel survives. If
the whole machine crashes, the uncommitted kernel transaction simply vanishes,
which is fine. The dangerous case is when the OSD process dies (or hits
ENOSPC) while the kernel keeps running and commits the half-done transaction.
A logical write-ahead log in user space is an alternative to using the file system's in-kernel transaction framework. It provides atomicity for transactions by writing down the full plan before doing any real work. Once the whole plan is safely on the disk, it adds one tiny "commit" record at the end. The disk writes this commit record completely or not at all.
If the system crashes, recovery just checks for the commit record. If it's there, recovery finishes the job. If it isn't, recovery ignores the plan, and since no real work had started, nothing needs undoing. Either way, we never get stuck with a half-done change.
11. Write the plan 2. Write COMMIT 3. Do the work2[ A-500, B+500 ] ───> [ COMMIT ✔ ] ───> [ change A and B ]3
4Crash? Check the log:5 COMMIT missing ──> ignore plan ──> nothing happened6 COMMIT present ──> redo plan ──> everything happenedIn terms of Ceph, the storage backend maintains its own WAL, called the journal. When a write comes in, the following happens:
fsync is called to commit the transaction to disk.1Transaction t = {2 write(obj A, 4 KiB @ 8192),3 setattr(A, v42),4 append(pglog, entry 42)5}6
7Step 1: Serialize t (including the 4 KiB of data) → append to journal8Step 2: fsync the journal ← COMMIT POINT. Now t is durable (safe).9Step 3: Apply t to XFS: pwrite(A), fsetxattr(A), pwrite(pglog)fsetxattr(A): when the replay
happens during recovery, we begin the entire transaction again from
pwrite(A).The journal gives us atomicity (one commit point) and durability. But it has three costs.
Cost 1: Slow read-modify-write
Ceph workloads usually follow a read-modify-write pattern. This is because the client often writes a smaller unit of data, but redundancy, checksums, etc. are all maintained over something bigger: a stripe, a blob, an omap entry set or an object version. Ceph has to bring that bigger unit up to date, and doing that requires knowing its current state. That is a read, then a modify, then a write.
The issue with the WAL is that a transaction's effect cannot be read by a
subsequent transaction until it has been committed to the journal (the fsync
step) and applied to the file system. This is because FileStore serves reads
from the file system, not from the journal.
1txn 1: increment counter in object X (10 → 11)2txn 2: increment counter in object X (needs to read 11)3
4time ──────────────────────────────────────────────►5txn 1: [serialize][journal write][fsync ~ms][apply to XFS]6txn 2: [read X = 11]7 [serialize][journal][fsync][apply]Every dependent operation now has to wait for the entire commit and apply before it can proceed.
Cost 2: Non-idempotent operations
Most operations done on the ObjectStore are idempotent; they are just set
operations (e.g., write(X, offset, data), truncate, setattr). There are a
few non-idempotent operations, such as clone(a → b), clone_range, remove
and split_collection. Each of these reads live state (data) and writes based
on it.
After a crash, replaying operations is only safe if each produces the same result when applied twice, but that's not the case with non-idempotent operations, since they depend on the live state of the object.
The flow of transactions is like this:
At any moment, recent transactions are in one of three states:
1txn 900–950 applied + synced → trimmable, gone after next sync2txn 951–990 applied, NOT synced → still in the journal (the risky window)3txn 991–1000 journaled, not applied yetsyncfs, the flush mechanism that moves data from the page cache
onto the physical diskWhen a transaction is "applied", the data usually lands in the Linux page cache in RAM, not on disk. The file system writes it out whenever it chooses. So "applied" doesn't mean "durable".
If the machine crashes, replay starts from the last sync point, i.e., txn 951. Transactions 951–990 might already be partly or fully on disk, and they might get applied again. That's where non-idempotent operations cause issues.
For example:
1txn:2 ① clone a→b3 ② update a4 ③ update cSay a crash occurs after ② has reached the disk. The physical storage now holds
the updated value of a. When replay starts again from ①, the clone copies the
updated value of a into b instead of the original value.
Cost 3: Double writes
The journal is a full write-ahead log. The whole transaction, including the data payload, is serialized and written to the journal and then applied to the file system. Every byte of data is written twice:
14 MiB write → 4 MiB into journal (fsync) → 4 MiB into XFSIn this approach, we put the transaction state into a key-value store rather than writing and managing our own log. This way, we hand off the hard part (the WAL machinery) to libraries like RocksDB, and the system just reads and writes keys.
In a hand-rolled WAL we need to design all of this:
fsyncWhen we use a KV store as the WAL, we can just write code like this:
1batch = new WriteBatch();2batch.put("obj/A/info", {version: 42, size: 12288});3batch.put("obj/A/attr/_", ...);4batch.put("pglog/42", "modify A v42");5db.write(batch, sync=true); // atomic + durable, doneThe KV store does exactly what we would have built. From the user's point of view, all we need to know is that once the write returns, the batch is durable and applied all-or-nothing.
The following are the conceptual changes:
Example: a tiny object store
1Keys in RocksDB:2 obj/photo1 → {data_at: extent 500–520, size: 81920, version: 3}3 obj/photo2 → {data_at: extent 900–902, size: 12000, version: 1}Writing a new photo:
11. allocate extents 1200–1240, write the data there, flush22. db.write({ put obj/photo3 → {data_at: 1200–1240, ...} }, sync)A crash before step 2 leaves extents 1200–1240 unreferenced: free space that a cleanup pass reclaims. A crash after step 2 means the photo fully exists.
Cost: Journal on top of a journal
The KV store keeps a diary (log) so that it can recover after a crash. But this KV store runs on top of a file system, and the file system also keeps its own diary to protect its files. So every time the KV store says "save my diary entry", two diaries get updated. These disk flushes are slow, often milliseconds on a hard drive, and having two flushes cuts how many objects we can create per second.
When the object's data also lives in a file on the same file system, creating one object looks like this:
1fsync(object's data file)2 ├─ flush the object's data ← flush 13 └─ flush XFS's own journal ← flush 2 (records the new file and its size)4fsync(RocksDB's log file)5 ├─ flush RocksDB's log entry ← flush 36 └─ flush XFS's own journal ← flush 4 (records "this log file grew")7─────────────────────────────────8total: 4 flushes (on a raw disk: 2)To summarize the three ways of building transactions on top of a file system:
| Approach | Atomicity comes from | What it costs |
|---|---|---|
| Borrow the kernel's transactions | The file system's internal journal | No rollback: a dead process leaves half a transaction committed |
| User-space WAL | Our own journal + fsync | Slow read-modify-write, unsafe replay, every byte written twice |
| KV store as WAL | A RocksDB write batch | Journal on a journal: 4 flushes where a raw disk needs 2 |
Metadata is information about the data that is stored. It contains everything we need to know about a piece of data in order to find it, manage it and trust it.
For a single 4 MB RADOS object, we have the following metadata:
rbd_data.1f2a.00032, and the PG and collection it
belongs toThis metadata may run to a few hundred bytes. Because metadata is what keeps the data findable and trustworthy, every data operation also updates it; a single data operation might require 4 or more metadata updates.
In distributed storage systems, where a lot of operations are performed on data, inefficient metadata operations become a bottleneck.
The objects on the file system would be stored as shown below:
1/var/lib/ceph/osd/ceph-0/current/2└── 2.3f_head/ ← one directory = one PG3 ├── rbd\udata.abc...09__head_0C11B03F__24 ├── rbd\udata.abc...05__head_E5326A3F__25 ├── rbd\udata.def...01__head_7A40913F__26 ├── rbd\udata.ghi...02__head_51C0D23F__27 ├── ...8 └── ... 1,000,000 files in the same directoryAnd the problem with storing metadata in the file system can be captured by the diagram below:
1/var/lib/ceph/osd/ceph-0/current/2.3f_head/2┌──────────────────────────────────────────────────────────────┐3│ rbd\udata.abc...09__head_0C11B03F__2 → inode 8812 ──┐ │4│ rbd\udata.abc...05__head_E5326A3F__2 → inode 1201 │ │5│ rbd\udata.def...01__head_7A40913F__2 → inode 4410 │ │6│ rbd\udata.ghi...02__head_51C0D23F__2 → inode 7733 │ │7│ ... 1,000,000 more entries ... │ │8└─────────────────────────────────────────────────────┼────────┘9 │ ▼10 │ inode 8812 (metadata)11 │ ┌──────────────────────────────┐12 │ │ size: 4 MB │13 │ │ owner, permissions, times │14 │ │ xattrs: version, attrs, │15 │ │ (real name if too long) │16 │ │ data: blocks 9000–10023 ─────┼──► object bytes17 │ └──────────────────────────────┘18 │19 │ ① readdir(): returns ALL 1,000,000 entries,20 │ in the file system's own order (random to Ceph)21 ▼22┌──────────────────────────────────────────┐23│ ..._0C11B03F, ..._E5326A3F, ..._7A40913F │ ← unsorted24└──────────────────────────────────────────┘25 │26 │ ② long names? read each inode's xattrs27 │ to get the real object name (1 inode lookup per entry)28 ▼29┌──────────────────────────────────────────┐30│ sort 1,000,000 names by hash │31└──────────────────────────────────────────┘32 │33 ▼34┌──────────────────────────────────────────┐35│ objects in hash order │ ← what Ceph actually needed36└──────────────────────────────────────────┘Enumerating objects is a very important operation for a distributed system, both for keeping data consistent and for repairs.
For Ceph, enumeration means "listing all objects that exist in a placement group". Ceph has no central metadata server holding the list of all objects. Each OSD only knows what it holds. So whenever the system needs answers to:
the only way to answer is to ask each OSD to list what it has and compare the lists.
Notice that one directory per PG already gives us only that PG's objects. The real problem is the order. Scrub and backfill don't ask for the whole list at once; they walk it in chunks, by hash:
1scrub: "give me the next 100 objects after hash 0x3A000000"2 → compare with the replica's next 100 → repeatHash order lets two replicas walk their lists in lockstep and compare them chunk
by chunk. It also lets a scrub pause and resume, locking only the range it is
currently checking. readdir() can't do this, since it returns the entries in
the file system's own order. So answering "the next 100 objects after H" means
reading and sorting the entire directory every single time, as shown in the
diagram above.
Newer drives like SMR HDDs and ZNS SSDs don't allow random overwrites; the software must write sequentially in zones. A backend built on a file system can't use them until the file system learns to, and that can take years. I am skipping this one in depth for this blog, since I need to understand it better myself.
A distributed system can only be as correct as its backends. Ceph needs each OSD's backend to do three things well:
File systems don't give us any of these cheaply. We have to fake transactions on top of them, the metadata ends up spread across directories, inodes and xattrs that can't be listed in the order Ceph needs, and new hardware has to wait until the file system supports it.
In Part 2, we will see Ceph run into exactly these problems, one backend at a time, from EBOFS to NewStore.