Genesis of Kosmos File System (KFS)

distributed systems
storage
Indexing tens of billions of URLs on commodity hardware in 2004 meant separating where data lived from where computation ran. That is why we built the Kosmos File System.
Author

Srinivasan Seshadri

Published

November 21, 2021

This post is about our experience building a web scale search engine, Kosmix, and, in particular how we felt the need for a distributed file system and built Kosmos File System. We started building the web search engine in 2004. We wanted to index as a large a corpus as we could and comparable web scale indexes were in the tens of billions of URLs then and so was ours.

We investigated Lucene (written in Java) as a starting point but soon realized Lucene would not be able to give us the query performance we needed on the index sizes we wanted on commodity hardware. We were buying commodity machines and putting them up in a data center. We decided to build our own web search engine in C/C++ and tune every component so we could get the performance we needed on this commodity hardware.

We had our own crawler, indexer, web graph based signal computation engine (we had our own categorization algorithm based on the web graph) and a real-time search engine. Our crawler would run continuously but to get the data ready for our real-time search engine we needed some expensive and long computations including building the inverted index and the web graph based signals. Building the inverted index was essentially a huge external distributed sort on tens of machines (we had tens of machines when we began and the aggregate memory of this system was tiny compared to the data we were indexing).

Building this inverted index took several weeks if there were no interruptions. Our first implementation was just one monolithic, long running computation and if something failed, we had to start all over again. Further, we had at least one machine failure every few weeks as we pushed these machines and kept them busy all the time — the disks were spinning continuously. A machine failure meant we had to take a trip to the data center to go fix it — sometimes needed a hard reboot to get things back online. This meant we needed two to three attempts to build the index before we could successfully complete. This meant on an average we took three months to update our user facing index. This was a completely unacceptable user experience.

It became clear to us that we needed a way of handling machine failures gracefully and given our computation was heavily dependent on data stored on disks, we needed a simple mechanism to both have multiple copies of the data as well as be able to access this data from any machine and run any computation on this data. Further, we had to break our single computation into small chunks of computation — this meant we would lose only that chunk of work if there was a failure. Thus was born the inspiration for Kosmos File System (KFS) (the GFS paper was published around that time and this reinforced our belief that we needed something like KFS).

Using database terminology of 80s and 90s, we moved from a shared nothing architecture to a shared disk (distributed file system) architecture and that allowed us to separate the compute from the storage. This resulted in huge flexibility in how we balanced the load across machines and ensure we were able to complete stragglers by pressing other machines into service. Large databases in the 80s and 90s were not built on top of commodity hardware and shared nothing architectures were easier to scale then. But web scale was a different league from the largest commercial databases and data warehouses and to make this economical scaling commodity hardware was the way to go — however coupled with the failure rates of commodity hardware this meant we needed a distributed file system or raw storage in the cloud (the word cloud was perhaps first used in 2006 ).