Skip to content
Naveen Naidu

Why Is BlueStore the Way It Is? Part 1: Why File Systems Fall Short

— 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:

  1. Why use a bespoke solution like BlueStore? Why are we not using already existing solutions like file systems?
  2. What were the previous solutions that were considered before we zeroed in on the current architecture?

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:

  1. Part 1 (this post): What a distributed system needs from its storage backend, and why file systems struggle to provide it.
  2. Part 2: Ceph's backends in order (EBOFS → FileStore on Btrfs → FileStore on XFS → NewStore) and how each of them ran into these problems.
  3. Part 3: How all of these learnings shaped BlueStore, and the price it pays for owning the disk.

Intro

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 C
6 │ │ │
7 BlueStore #0 BlueStore #1 BlueStore #2 ← "distributed storage backends"
8 │ │ │ (three independent, local instances)
9 /dev/sdb /dev/sdb /dev/sdb

What Does a Distributed System Need from Its Backend?

In 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:

  • The object's data
  • The object's metadata
  • An entry in the placement group's log

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:

  1. Strong consistency, which needs atomicity, durability, ordering and integrity.
  2. Efficiency, which needs cheap transactions and fast metadata.
  3. Longevity, which needs hardware flexibility.

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 needsWhich is hard on a file system because of
Strong consistencyTransactions: there are no multi-op atomic transactions
EfficiencyTransactions (costly to fake) and Metadata (slow with millions of objects)
LongevityNew hardware: we have to wait for the file system to support it

File Systems as Storage Backends

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:

  • It started with EBOFS, a custom user-space object store (added in 2005).
  • It moved to FileStore on Btrfs (~2008) to get transactions, checksums and compression.
  • Then XFS, for stability.
  • Then back to the raw disk with BlueStore.

We'll explore each of them in detail in Part 2.

File systems became the de facto standard for the following reasons:

  • Delegate the hard parts. Data persistence, crash-safe metadata, allocation and caching come from well-tested, mature and highly performant code. A team building a distributed system need not reinvent the wheel.
  • A familiar interface. Files, directories and POSIX are easy to reason about.
  • Tooling. Standard tools such as ls and find can be used to explore the disk contents.
  • Two environmental reasons:
    • Linux's ubiquity.
    • The workload fit: the workloads of that time did not require millions of small objects, overwrites or multi-op atomicity.

Problems with File Systems as Storage Backends

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:

  1. Inefficient transactions
  2. Slow metadata operations
  3. Less support for new storage hardware

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.

Problem 1: Transactions

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 file
2int fd = open("config.tmp", O_WRONLY | O_CREAT | O_TRUNC, 0644);
3write(fd, new_contents, len);
4
5// 2. Make the new contents durable
6fsync(fd);
7close(fd);
8
9// 3. Atomically swap it into place
10rename("config.tmp", "config");
11
12// 4. Make the swap itself durable
13int 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)
2fsync
3op2: set xattr version = 42
4fsync
5op3: append PG log entry 42
6fsync

If 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:

  1. Hooking into the file system's internal transaction mechanism (FileStore on Btrfs)
  2. Implementing a write-ahead log (WAL) in user space (FileStore on XFS)
  3. Using a key-value database with transactions as a WAL (NewStore)

The next sections talk about why these options have significant performance/complexity overhead.

Leveraging the File System's Internal Transactions

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 space
2 // modify bitmap block in memory
3 // modify inode block in memory
4 // modify directory block in memory
5jbd2_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:

  1. Not exposed. Applications can't normally use them at all.
  2. No rollback. Kernel file operations are designed to never fail halfway. Before modifying anything, the kernel does all the checks that could fail (permissions, free space, journal credits, etc.) to prevent mid-operation failures.
  3. No clear boundaries. Some file systems batch many operations from many processes into one big transaction every few seconds. There is no notion of segregation, i.e., nothing that says "these three ops from this process, and only these".

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_START
2 op1 ✅ applied
3 ← process dies
4 op2 never issued
5 op3 never issued
6TRANS_END never issued

Since 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.

Implementing the WAL in User Space

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 work
2[ A-500, B+500 ] ───> [ COMMIT ✔ ] ───> [ change A and B ]
3
4Crash? Check the log:
5 COMMIT missing ──> ignore plan ──> nothing happened
6 COMMIT present ──> redo plan ──> everything happened

In terms of Ceph, the storage backend maintains its own WAL, called the journal. When a write comes in, the following happens:

  1. The transaction is serialized and written to the journal.
  2. fsync is called to commit the transaction to disk.
  3. The operations in the transaction are applied to the 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 journal
8Step 2: fsync the journal ← COMMIT POINT. Now t is durable (safe).
9Step 3: Apply t to XFS: pwrite(A), fsetxattr(A), pwrite(pglog)
  • Crash before step 2: the journal entry is incomplete, so ignore it. Treat it as if the transaction never happened.
  • Crash after step 2: on restart, we replay all the journal entries and they get reapplied. Say the crash happens at 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:

  1. Write the transaction into the journal.
  2. Apply it to the file system.
  3. Some time later, the journal is trimmed and the old entries are thrown away.

At any moment, recent transactions are in one of three states:

1txn 900–950 applied + synced → trimmable, gone after next sync
2txn 951–990 applied, NOT synced → still in the journal (the risky window)
3txn 991–1000 journaled, not applied yet
  • applied = data has been written into the page cache (RAM)
  • synced = syncfs, the flush mechanism that moves data from the page cache onto the physical disk

When 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→b
3 ② update a
4 ③ update c

Say 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:

  1. once into the journal
  2. later into the file system
14 MiB write → 4 MiB into journal (fsync) → 4 MiB into XFS

Using a Key-Value Store as the WAL

In 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:

  • a log file format
  • a method to serialize transactions into it
  • commit records + fsync
  • replay logic after a crash (this must be idempotent)

When 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, done

The 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:

  1. We now store state, not a log of operations. In a WAL, we recorded operations that were applied elsewhere; in a KV store, we write the resulting state to the associated key.
  2. The KV store becomes the source of truth. Since we have keys, it's easy to answer questions like "Does object A exist?", "Where's its data?", "What version is it?" You can just look up the key instead of inferring the answer from file names and file sizes.

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, flush
22. 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 1
3 └─ 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 3
6 └─ flush XFS's own journal ← flush 4 (records "this log file grew")
7─────────────────────────────────
8total: 4 flushes (on a raw disk: 2)

Summary of the three approaches

To summarize the three ways of building transactions on top of a file system:

ApproachAtomicity comes fromWhat it costs
Borrow the kernel's transactionsThe file system's internal journalNo rollback: a dead process leaves half a transaction committed
User-space WALOur own journal + fsyncSlow read-modify-write, unsafe replay, every byte written twice
KV store as WALA RocksDB write batchJournal on a journal: 4 flushes where a raw disk needs 2

Problem 2: Metadata Operations

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:

  1. Object name, e.g. rbd_data.1f2a.00032, and the PG and collection it belongs to
  2. Attributes: size, version, snapshot, info, xattrs, omap keys
  3. Location: the disk extents that hold its bytes, which OSD
  4. Integrity and history: checksums, PG log entry

This 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 PG
3 ├── rbd\udata.abc...09__head_0C11B03F__2
4 ├── rbd\udata.abc...05__head_E5326A3F__2
5 ├── rbd\udata.def...01__head_7A40913F__2
6 ├── rbd\udata.ghi...02__head_51C0D23F__2
7 ├── ...
8 └── ... 1,000,000 files in the same directory

And 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 bytes
17 │ └──────────────────────────────┘
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 │ ← unsorted
24└──────────────────────────────────────────┘
25 │
26 │ ② long names? read each inode's xattrs
27 │ 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 needed
36└──────────────────────────────────────────┘

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:

  • Which objects is this replica missing?
  • Do all replicas hold the same objects with the same contents?

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 → repeat

Hash 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.

Problem 3: New Storage Hardware

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.

Wrapping Up

A distributed system can only be as correct as its backends. Ceph needs each OSD's backend to do three things well:

  1. Apply multi-op transactions atomically and cheaply
  2. Handle metadata fast, including listing millions of objects in order
  3. Adopt new storage hardware without waiting on someone else

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.