Wednesday, October 13, 2021

GFS: Evolution on Fast-forward | 20 minutes to read | Understand Google file system design

Oct. 13, 2021

Here is the article.

GFS: Evolution on Fast-forward

A discussion between Kirk McKusick and Sean Quinlan about the origin and evolution of the Google File System.

During the early stages of development at Google, the initial thinking did not include plans for building a new file system. While work was still being done on one of the earliest versions of the company's crawl and indexing system, however, it became quite clear to the core engineers that they really had no other choice, and GFS (Google File System) was born.

First, given that Google's goal was to build a vast storage network out of inexpensive commodity hardware, it had to be assumed that component failures would be the norm—meaning that constant monitoring, error detection, fault tolerance, and automatic recovery would have to be an integral part of the file system. Also, even by Google's earliest estimates, the system's throughput requirements were going to be daunting by anybody's standards—featuring multi-gigabyte files and data sets containing terabytes of information and millions of objects. Clearly, this meant traditional assumptions about I/O operations and block sizes would have to be revisited. There was also the matter of scalability. This was a file system that would surely need to scale like no other. Of course, back in those earliest days, no one could have possibly imagined just how much scalability would be required. They would learn about that soon enough.

Still, nearly a decade later, most of Google's mind-boggling store of data and its ever-growing array of applications continue to rely upon GFS. Many adjustments have been made to the file system along the way, and—together with a fair number of accommodations implemented within the applications that use GFS—they have made the journey possible.

To explore the reasoning behind a few of the more crucial initial design decisions as well as some of the incremental adaptations that have been made since then, ACM asked Sean Quinlan to pull back the covers on the changing file-system requirements and the evolving thinking at Google. Since Quinlan served as the GFS tech leader for a couple of years and continues now as a principal engineer at Google, he's in a good position to offer that perspective. As a grounding point beyond the Googleplex, ACM asked Kirk McKusick to lead the discussion. He is best known for his work on BSD (Berkeley Software Distribution) Unix, including the original design of the Berkeley FFS (Fast File System).

The discussion starts, appropriately enough, at the beginning—with the unorthodox decision to base the initial GFS implementation on a single-master design. At first blush, the risk of a single centralized master becoming a bandwidth bottleneck—or, worse, a single point of failure—seems fairly obvious, but it turns out Google's engineers had their reasons for making this choice.


MCKUSICK One of the more interesting—and significant—aspects of the original GFS architecture was the decision to base it on a single master. Can you walk us through what led to that decision?

QUINLAN The decision to go with a single master was actually one of the very first decisions, mostly just to simplify the overall design problem. That is, building a distributed master right from the outset was deemed too difficult and would take too much time. Also, by going with the single-master approach, the engineers were able to simplify a lot of problems. Having a central place to control replication and garbage collection and many other activities was definitely simpler than handling it all on a distributed basis. So the decision was made to centralize that in one machine.

MCKUSICK Was this mostly about being able to roll out something within a reasonably short time frame?

QUINLAN Yes. In fact, some of the engineers who were involved in that early effort later went on to build BigTable, a distributed storage system, but that effort took many years. The decision to build the original GFS around the single master really helped get something out into the hands of users much more rapidly than would have otherwise been possible.

Also, in sketching out the use cases they anticipated, it didn't seem the single-master design would cause much of a problem. The scale they were thinking about back then was framed in terms of hundreds of terabytes and a few million files. In fact, the system worked just fine to start with.

MCKUSICK But then what?

QUINLAN Problems started to occur once the size of the underlying storage increased. Going from a few hundred terabytes up to petabytes, and then up to tens of petabytes that really required a proportionate increase in the amount of metadata the master had to maintain. Also, operations such as scanning the metadata to look for recoveries all scaled linearly with the volume of data. So the amount of work required of the master grew substantially. The amount of storage needed to retain all that information grew as well.

In addition, this proved to be a bottleneck for the clients, even though the clients issue few metadata operations themselves—for example, a client talks to the master whenever it does an open. When you have thousands of clients all talking to the master at the same time, given that the master is capable of doing only a few thousand operations a second, the average client isn't able to command all that many operations per second. Also bear in mind that there are applications such as MapReduce, where you might suddenly have a thousand tasks, each wanting to open a number of files. Obviously, it would take a long time to handle all those requests, and the master would be under a fair amount of duress.

MCKUSICK Now, under the current schema for GFS, you have one master per cell, right?

QUINLAN That's correct.

MCKUSICK And historically you've had one cell per data center, right?

QUINLAN That was initially the goal, but it didn't work out like that to a large extent—partly because of the limitations of the single-master design and partly because isolation proved to be difficult. As a consequence, people generally ended up with more than one cell per data center. We also ended up doing what we call a "multi-cell" approach, which basically made it possible to put multiple GFS masters on top of a pool of chunkservers. That way, the chunkservers could be configured to have, say, eight GFS masters assigned to them, and that would give you at least one pool of underlying storage—with multiple master heads on it, if you will. Then the application was responsible for partitioning data across those different cells.

MCKUSICK Presumably each application would then essentially have its own master that would be responsible for managing its own little file system. Was that basically the idea?

QUINLAN Well, yes and no. Applications would tend to use either one master or a small set of the masters. We also have something we called Name Spaces, which are just a very static way of partitioning a namespace that people can use to hide all of this from the actual application. The Logs Processing System offers an example of this approach: once logs exhaust their ability to use just one cell, they move to multiple GFS cells; a namespace file describes how the log data is partitioned across those different cells and basically serves to hide the exact partitioning from the application. But this is all fairly static.

MCKUSICK What's the performance like, in light of all that?

QUINLAN We ended up putting a fair amount of effort into tuning master performance, and it's atypical of Google to put a lot of work into tuning any one particular binary. Generally, our approach is just to get things working reasonably well and then turn our focus to scalability—which usually works well in that you can generally get your performance back by scaling things. Because in this instance we had a single bottleneck that was starting to have an impact on operations, however, we felt that investing a bit of additional effort into making the master lighter weight would be really worthwhile. In the course of scaling from thousands of operations to tens of thousands and beyond, the single master had become somewhat less of a bottleneck. That was a case where paying more attention to the efficiency of that one binary definitely helped keep GFS going for quite a bit longer than would have otherwise been possible.


It could be argued that managing to get GFS ready for production in record time constituted a victory in its own right and that, by speeding Google to market, this ultimately contributed mightily to the company's success. A team of three was responsible for all of that—for the core of GFS—and for the system being readied for deployment in less than a year.

But then came the price that so often befalls any successful system—that is, once the scale and use cases have had time to expand far beyond what anyone could have possibly imagined. In Google's case, those pressures proved to be particularly intense.

Although organizations don't make a habit of exchanging file-system statistics, it's safe to assume that GFS is the largest file system in operation (in fact, that was probably true even before Google's acquisition of YouTube). Hence, even though the original architects of GFS felt they had provided adequately for at least a couple of orders of magnitude of growth, Google quickly zoomed right past that.

In addition, the number of applications GFS was called upon to support soon ballooned. In an interview with one of the original GFS architects, Howard Gobioff (conducted just prior to his surprising death in early 2008), he recalled, "The original consumer of all our earliest GFS versions was basically this tremendously large crawling and indexing system. The second wave came when our quality team and research groups started using GFS rather aggressively—and basically, they were all looking to use GFS to store large data sets. And then, before long, we had 50 users, all of whom required a little support from time to time so they'd all keep playing nicely with each other."

One thing that helped tremendously was that Google built not only the file system but also all of the applications running on top of it. While adjustments were continually made in GFS to make it more accommodating to all the new use cases, the applications themselves were also developed with the various strengths and weaknesses of GFS in mind. "Because we built everything, we were free to cheat whenever we wanted to," Gobioff neatly summarized. "We could push problems back and forth between the application space and the file-system space, and then work out accommodations between the two."

The matter of sheer scale, however, called for some more substantial adjustments. One coping strategy had to do with the use of multiple "cells" across the network, functioning essentially as related but distinct file systems. Besides helping to deal with the immediate problem of scale, this proved to be a more efficient arrangement for the operations of widely dispersed data centers.

Rapid growth also put pressure on another key parameter of the original GFS design: the choice to establish 64 MB as the standard chunk size. That, of course, was much larger than the typical file-system block size, but only because the files generated by Google's crawling and indexing system were unusually large. As the application mix changed over time, however, ways had to be found to let the system deal efficiently with large numbers of files requiring far less than 64 MB (think in terms of Gmail, for example). The problem was not so much with the number of files itself, but rather with the memory demands all of those files made on the centralized master, thus exposing one of the bottleneck risks inherent in the original GFS design.


MCKUSICK I gather from the original GFS paper [Ghemawat, S., Gobioff, H.,  Leung, S-T. 2003. The Google File System. SOSP (ACM Symposium on Operating Systems Principles)] that file counts have been a significant issue for you right along. Can you go into that a little bit?                    

QUINLAN The file-count issue came up fairly early because of the way people ended up designing their systems around GFS. Let me cite a specific example. Early in my time at Google, I was involved in the design of the Logs Processing system. We initially had a model where a front-end server would write a log, which we would then basically copy into GFS for processing and archival. That was fine to start with, but then the number of front-end servers increased, each rolling logs every day. At the same time, the number of log types was going up, and then you'd have front-end servers that would go through crash loops and generate lots more logs. So we ended up with a lot more files than we had anticipated based on our initial back-of-the-envelope estimates.

This became an area we really had to keep an eye on. Finally, we just had to concede there was no way we were going to survive a continuation of the sort of file-count growth we had been experiencing.

MCKUSICK Let me make sure I'm following this correctly: your issue with file-count growth is a result of your needing to have a piece of metadata on the master for each file, and that metadata has to fit in the master's memory.

QUINLAN That's correct.

MCKUSICK And there are only a finite number of files you can accommodate before the master runs out of memory?

QUINLAN Exactly. And there are two bits of metadata. One identifies the file, and the other points out the chunks that back that file. If you had a chunk that contained only 1 MB, it would take up only 1 MB of disk space, but it still would require those two bits of metadata on the master. If your average file size ends up dipping below 64 MB, the ratio of the number of objects on your master to what you have in storage starts to go down. That's where you run into problems.

Going back to that logs example, it quickly became apparent that the natural mapping we had thought of—and which seemed to make perfect sense back when we were doing our back-of-the-envelope estimates—turned out not to be acceptable at all. We needed to find a way to work around this by figuring out how we could combine some number of underlying objects into larger files. In the case of the logs, that wasn't exactly rocket science, but it did require a lot of effort.

MCKUSICK That sounds like the old days when IBM had only a minimum disk allocation, so it provided you with a utility that let you pack a bunch of files together and then create a table of contents for that.

QUINLAN Exactly. For us, each application essentially ended up doing that to varying degrees. That proved to be less burdensome for some applications than for others. In the case of our logs, we hadn't really been planning to delete individual log files. It was more likely that we would end up rewriting the logs to anonymize them or do something else along those lines. That way, you don't get the garbage-collection problems that can come up if you delete only some of the files within a bundle.

For some other applications, however, the file-count problem was more acute. Many times, the most natural design for some application just wouldn't fit into GFS—even though at first glance you would think the file count would be perfectly acceptable, it would turn out to be a problem. When we started using more shared cells, we put quotas on both file counts and storage space. The limit that people have ended up running into most has been, by far, the file-count quota. In comparison, the underlying storage quota rarely proves to be a problem.

MCKUSICK What longer-term strategy have you come up with for dealing with the file-count issue? Certainly, it doesn't seem that a distributed master is really going to help with that—not if the master still has to keep all the metadata in memory, that is.

QUINLAN The distributed master certainly allows you to grow file counts, in line with the number of machines you're willing to throw at it. That certainly helps. 

One of the appeals of the distributed multimaster model is that if you scale everything up by two orders of magnitude, then getting down to a 1-MB average file size is going to be a lot different from having a 64-MB average file size. If you end up going below 1 MB, then you're also going to run into other issues that you really need to be careful about. For example, if you end up having to read 10,000 10-KB files, you're going to be doing a lot more seeking than if you're just reading 100 1-MB files.

My gut feeling is that if you design for an average 1-MB file size, then that should provide for a much larger class of things than does a design that assumes a 64-MB average file size. Ideally, you would like to imagine a system that goes all the way down to much smaller file sizes, but 1 MB seems a reasonable compromise in our environment.

MCKUSICK What have you been doing to design GFS to work with 1-MB files?

QUINLAN We haven't been doing anything with the existing GFS design. Our distributed master system that will provide for 1-MB files is essentially a whole new design. That way, we can aim for something on the order of 100 million files per master. You can also have hundreds of masters.

MCKUSICK So, essentially no single master would have all this data on it?

QUINLAN That's the idea.


With the recent emergence within Google of BigTable, a distributed storage system for managing structured data, one potential remedy for the file-count problem—albeit perhaps not the very best one—is now available.

The significance of BigTable goes far beyond file counts, however. Specifically, it was designed to scale into the petabyte range across hundreds or thousands of machines, as well as to make it easy to add more machines to the system and automatically start taking advantage of those resources without reconfiguration. For a company predicated on the notion of employing the collective power, potential redundancy, and economies of scale inherent in a massive deployment of commodity hardware, these rate as significant advantages indeed.

Accordingly, BigTable is now used in conjunction with a growing number of Google applications. Although it represents a departure of sorts from the past, it also must be said that BigTable was built on GFS, runs on GFS, and was consciously designed to remain consistent with most GFS principles. Consider it, therefore, as one of the major adaptations made along the way to help keep GFS viable in the face of rapid and widespread change.


MCKUSICK You now have this thing called BigTable. Do you view that as an application in its own right?

QUINLAN From the GFS point of view, it's an application, but it's clearly more of an infrastructure piece.

MCKUSICK If I understand this correctly, BigTable is essentially a lightweight relational database.

QUINLAN It's not really a relational database. I mean, we're not doing SQL and it doesn't really support joins and such. But BigTable is a structured storage system that lets you have lots of key-value pairs and a schema.

MCKUSICK Who are the real clients of BigTable?

QUINLAN BigTable is increasingly being used within Google for crawling and indexing systems, and we use it a lot within many of our client-facing applications. The truth of the matter is that there are tons of BigTable clients. Basically, any app with lots of small data items tends to use BigTable. That's especially true wherever there's fairly structured data.

MCKUSICK I guess the question I'm really trying to pose here is: Did BigTable just get stuck into a lot of these applications as an attempt to deal with the small-file problem, basically by taking a whole bunch of small things and then aggregating them together?

QUINLAN That has certainly been one use case for BigTable, but it was actually intended for a much more general sort of problem. If you're using BigTable in that way—that is, as a way of fighting the file-count problem where you might have otherwise used a file system to handle that—then you would not end up employing all of BigTable's functionality by any means. BigTable isn't really ideal for that purpose in that it requires resources for its own operations that are nontrivial. Also, it has a garbage-collection policy that's not super-aggressive, so that might not be the most efficient way to use your space. I'd say that the people who have been using BigTable purely to deal with the file-count problem probably haven't been terribly happy, but there's no question that it is one way for people to handle that problem.

MCKUSICK What I've read about GFS seems to suggest that the idea was to have only two basic data structures: logs and SSTables (Sorted String Tables). Since I'm guessing the SSTables must be used to handle key-value pairs and that sort of thing, how is that different from BigTable?

QUINLAN The main difference is that SSTables are immutable, while BigTable provides mutable key value storage, and a whole lot more. BigTable itself is actually built on top of logs and SSTables. Initially, it stores incoming data into transaction log files. Then it gets compacted—as we call it—into a series of SSTables, which in turn get compacted together over time. In some respects, it's reminiscent of a log-structure file system. Anyway, as you've observed, logs and SSTables do seem to be the two data structures underlying the way we structure most of our data. We have log files for mutable stuff as it's being recorded. Then, once you have enough of that, you sort it and put it into this structure that has an index.

Even though GFS does not provide a Posix interface, it still has a pretty generic file-system interface, so people are essentially free to write any sort of data they like. It's just that, over time, the majority of our users have ended up using these two data structures. We also have something called protocol buffers, which is our data description language. The majority of data ends up being protocol buffers in these two structures.

Both provide for compression and checksums. Even though there are some people internally who end up reinventing these things, most people are content just to use those two basic building blocks.


Because GFS was designed initially to enable a crawling and indexing system, throughput was everything. In fact, the original paper written about the system makes this quite explicit: "High sustained bandwidth is more important than low latency. Most of our target applications place a premium on processing data in bulk at a high rate, while few have stringent response-time requirements for an individual read and write."

But then Google either developed or embraced many user-facing Internet services for which this is most definitely not the case.

One GFS shortcoming that this immediately exposed had to do with the original single-master design. A single point of failure may not have been a disaster for batch-oriented applications, but it was certainly unacceptable for latency-sensitive applications, such as video serving. The later addition of automated failover capabilities helped, but even then service could be out for up to a minute.

The other major challenge for GFS, of course, has revolved around finding ways to build latency-sensitive applications on top of a file system designed around an entirely different set of priorities.  


MCKUSICK It's well documented that the initial emphasis in designing GFS was on batch efficiency as opposed to low latency. Now that has come back to cause you trouble, particularly in terms of handling things such as videos. How are you handling that?

QUINLAN The GFS design model from the get-go was all about achieving throughput, not about the latency at which that might be achieved. To give you a concrete example, if you're writing a file, it will typically be written in triplicate—meaning you'll actually be writing to three chunkservers. Should one of those chunkservers die or hiccup for a long period of time, the GFS master will notice the problem and schedule what we call a pullchunk, which means it will basically replicate one of those chunks. That will get you back up to three copies, and then the system will pass control back to the client, which will continue writing.

When we do a pullchunk we limit it to something on the order of 5-10 MB a second. So, for 64 MB, you're talking about 10 seconds for this recovery to take place. There are lots of other things like this that might take 10 seconds to a minute, which works just fine for batch-type operations. If you're doing a large MapReduce operation, you're OK just so long as one of the items is not a real straggler, in which case you've got yourself a different sort of problem. Still, generally speaking, a hiccup on the order of a minute over the course of an hour-long batch job doesn't really show up. If you are working on Gmail, however, and you're trying to write a mutation that represents some user action, then getting stuck for a minute is really going to mess you up.

We've had similar issues with our master failover. Initially, GFS had no provision for automatic master failover. It was a manual process. Although it didn't happen a lot, whenever it did, the cell might be down for an hour. Even our initial master-failover implementation required on the order of minutes. Over the past year, however, we've taken that down to something on the order of tens of seconds.

MCKUSICK Still, for user-facing applications, that's not acceptable.

QUINLAN Right. While these instances—where you have to provide for failover and error recovery—may have been acceptable in the batch situation, they're definitely not OK from a latency point of view for a user-facing application. Another issue here is that there are places in the design where we've tried to optimize for throughput by dumping thousands of operations into a queue and then just processing through them. That leads to fine throughput, but it's not great for latency. You can easily get into situations where you might be stuck for seconds at a time in a queue just waiting to get to the head of the queue.

Our user base has definitely migrated from being a MapReduce-based world to more of an interactive world that relies on things such as BigTable. Gmail is an obvious example of that. Videos aren't quite as bad where GFS is concerned because you get to stream data, meaning you can buffer. Still, trying to build an interactive database on top of a file system that was designed from the start to support more batch-oriented operations has certainly proved to be a pain point.

MCKUSICK How exactly have you managed to deal with that?

QUINLAN Within GFS, we've managed to improve things to a certain degree, mostly by designing the applications to deal with the problems that come up. Take BigTable as a good concrete example. The BigTable transaction log is actually the biggest bottleneck for getting a transaction logged. In effect, we decided, "Well, we're going to see hiccups in these writes, so what we'll do is to have two logs open at any one time. Then we'll just basically merge the two. We'll write to one and if that gets stuck, we'll write to the other. We'll merge those logs once we do a replay—if we need to do a replay, that is." We tended to design our applications to function like that—which is to say they basically try to hide that latency since they know the system underneath isn't really all that great.

The guys who built Gmail went to a multihomed model, so if one instance of your Gmail account got stuck, you would basically just get moved to another data center. Actually, that capability was needed anyway just to ensure availability. Still, part of the motivation was that they wanted to hide the GFS problems.

MCKUSICK I think it's fair to say that, by moving to a distributed-master file system, you're definitely going to be able to attack some of those latency issues.

QUINLAN That was certainly one of our design goals. Also, BigTable itself is a very failure-aware system that tries to respond to failures far more rapidly than we were able to before. Using that as our metadata storage helps with some of those latency issues as well.


The engineers who worked on the earliest versions of GFS weren't particularly shy about departing from traditional choices in file-system design whenever they felt the need to do so. It just so happens that the approach taken to consistency is one of the aspects of the system where this is particularly evident.

Part of this, of course, was driven by necessity. Since Google's plans rested largely on massive deployments of commodity hardware, failures and hardware-related faults were a given. Beyond that, according to the original GFS paper, there were a few compatibility issues. "Many of our disks claimed to the Linux driver that they supported a range of IDE protocol versions but in fact responded reliably only to the more recent ones. Since the protocol versions are very similar, these drives mostly worked but occasionally the mismatches would cause the drive and the kernel to disagree about the drive's state. This would corrupt data silently due to problems in the kernel. This problem motivated our use of checksums to detect data corruption."

That didn't mean just any checksumming, however, but instead rigorous end-to-end checksumming, with an eye to everything from disk corruption to TCP/IP corruption to machine backplane corruption.

Interestingly, for all that checksumming vigilance, the GFS engineering team also opted for an approach to consistency that's relatively loose by file-system standards. Basically, GFS simply accepts that there will be times when people will end up reading slightly stale data. Since GFS is used mostly as an append-only system as opposed to an overwriting system, this generally means those people might end up missing something that was appended to the end of the file after they'd already opened it. To the GFS designers, this seemed an acceptable cost (although it turns out that there are applications for which this proves problematic).

Also, as Gobioff explained, "The risk of stale data in certain circumstances is just inherent to a highly distributed architecture that doesn't ask the master to maintain all that much information. We definitely could have made things a lot tighter if we were willing to dump a lot more data into the master and then have it maintain more state. But that just really wasn't all that critical to us."

Perhaps an even more important issue here is that the engineers making this decision owned not just the file system but also the applications intended to run on the file system. According to Gobioff, "The thing is that we controlled both the horizontal and the vertical—the file system and the application. So we could be sure our applications would know what to expect from the file system. And we just decided to push some of the complexity out to the applications to let them deal with it."

Still, there are some at Google who wonder whether that was the right call if only because people can sometimes obtain different data in the course of reading a given file multiple times, which tends to be so strongly at odds with their whole notion of how data storage is supposed to work.


MCKUSICK Let's talk about consistency. The issue seems to be that it presumably takes some amount of time to get everything fully written to all the replicas. I think you said something earlier to the effect that GFS essentially requires that this all be fully written before you can continue.

QUINLAN That's correct.

MCKUSICK If that's the case, then how can you possibly end up with things that aren't consistent?

QUINLAN Client failures have a way of fouling things up. Basically, the model in GFS is that the client just continues to push the write until it succeeds. If the client ends up crashing in the middle of an operation, things are left in a bit of an indeterminate state.

Early on, that was sort of considered to be OK, but over time, we tightened the window for how long that inconsistency could be tolerated, and then we slowly continued to reduce that. Otherwise, whenever the data is in that inconsistent state, you may get different lengths for the file. That can lead to some confusion. We had to have some backdoor interfaces for checking the consistency of the file data in those instances. We also have something called RecordAppend, which is an interface designed for multiple writers to append to a log concurrently. There the consistency was designed to be very loose. In retrospect, that turned out to be a lot more painful than anyone expected.

MCKUSICK What exactly was loose? If the primary replica picks what the offset is for each write and then makes sure that actually occurs, I don't see where the inconsistencies are going to come up.

QUINLAN What happens is that the primary will try. It will pick an offset, it will do the writes, but then one of them won't actually get written. Then the primary might change, at which point it can pick a different offset. RecordAppend does not offer any replay protection either. You could end up getting the data multiple times in the file.

There were even situations where you could get the data in a different order. It might appear multiple times in one chunk replica, but not necessarily in all of them. If you were reading the file, you could discover the data in different ways at different times. At the record level, you could discover the records in different orders depending on which chunks you happened to be reading.

MCKUSICK Was this done by design?

QUINLAN At the time, it must have seemed like a good idea, but in retrospect I think the consensus is that it proved to be more painful than it was worth. It just doesn't meet the expectations people have of a file system, so they end up getting surprised. Then they had to figure out work-arounds.

MCKUSICK In retrospect, how would you handle this differently?

QUINLAN I think it makes more sense to have a single writer per file.

MCKUSICK All right, but what happens when you have multiple people wanting to append to a log?

QUINLAN You serialize the writes through a single process that can ensure the replicas are consistent.

MCKUSICK There's also this business where you essentially snapshot a chunk. Presumably, that's something you use when you're essentially replacing a replica, or whenever some chunkserver goes down and you need to replace some of its files.

QUINLAN Actually, two things are going on there. One, as you suggest, is the recovery mechanism, which definitely involves copying around replicas of the file. The way that works in GFS is that we basically revoke the lock so that the client can't write it anymore, and this is part of that latency issue we were talking about.

There's also a separate issue, which is to support the snapshot feature of GFS. GFS has the most general-purpose snapshot capability you can imagine. You could snapshot any directory somewhere, and then both copies would be entirely equivalent. They would share the unchanged data. You could change either one and you could further snapshot either one. So it was really more of a clone than what most people think of as a snapshot. It's an interesting thing, but it makes for difficulties—especially as you try to build more distributed systems and you want potentially to snapshot larger chunks of the file tree.

I also think it's interesting that the snapshot feature hasn't been used more since it's actually a very powerful feature. That is, from a file-system point of view, it really offers a pretty nice piece of functionality. But putting snapshots into file systems, as I'm sure you know, is a real pain.

MCKUSICK:  I know. I've done it. It's excruciating—especially in an overwriting file system.

QUINLAN Exactly. This is a case where we didn't cheat, but from an implementation perspective, it's hard to create true snapshots. Still, it seems that in this case, going the full deal was the right decision. Just the same, it's an interesting contrast to some of the other decisions that were made early on in terms of the semantics.


All in all, the report card on GFS nearly 10 years later seems positive. There have been problems and shortcomings, to be sure, but there's surely no arguing with Google's success and GFS has without a doubt played an important role in that. What's more, its staying power has been nothing short of remarkable given that Google's operations have scaled orders of magnitude beyond anything the system had been designed to handle, while the application mix Google currently supports is not one that anyone could have possibly imagined back in the late '90s.

Still, there's no question that GFS faces many challenges now. For one thing, the awkwardness of supporting an ever-growing fleet of user-facing, latency-sensitive applications on top of a system initially designed for batch-system throughput is something that's obvious to all.

The advent of BigTable has helped somewhat in this regard. As it turns out, however, BigTable isn't actually all that great a fit for GFS. In fact, it just makes the bottleneck limitations of the system's single-master design more apparent than would otherwise be the case.

For these and other reasons, engineers at Google have been working for much of the past two years on a new distributed master system designed to take full advantage of BigTable to attack some of those problems that have proved particularly difficult for GFS.

Accordingly, it now seems that beyond all the adjustments made to ensure the continued survival of GFS, the newest branch on the evolutionary tree will continue to grow in significance over the years to come.
Q

LOVE IT, HATE IT? LET US KNOW

feedback@queue.acm.org

© 2009 ACM 1542-7730/09/0800 $10.00

acmqueue

Originally published in Queue vol. 7, no. 7—
see this item in the ACM Digital Library

GFS paper | Chinese article

Oct. 13, 2021

Here is the article. 

 下面是本课程要阅读的第二篇论文,同样鼎鼎大名的GFS(其实中间还有一课是对Go语言和RPC的讲解,在此就略过了。)Again,本人水平有限,请谨慎参考。

GFS为上一篇介绍的MapReduce提供了存储,没有GFS也就没有MapReduce了。其实与其说MapReduce多么牛,不如说是GFS牛。MapReduce论文出来的时候被做数据库的喷,不是没有道理的,这个模型早就是数据库领域几十年前玩剩下的了。但他们为什么做不出Google里面那种廉价高效的系统呢,我觉得主要是因为下层的支持系统不好做。

应用场景

GFS是一个分布式的文件系统。它不追求与POSIX接口兼容(我觉得这里不强求兼容是对的)而是使用自己的一套API。它对主要应用场景的假设也与传统单机文件系统不同:

  • 系统用大量普通机器堆起来,机器挂掉是非常常见的事。
  • 系统里面主要存的是大文件,几百M甚至上G这种大小。
  • 读操作是大量的顺序读和少量的随机读。
  • 写操作是大量的顺序写(append),一旦写入几乎不会修改。
  • 文件经常被大量client同时append,需要仔细考虑同步机制的开销
  • 相比低latency,更看重高吞吐量

可以看出这些需求基本就是针对MapReduce量身订做的。

系统架构

系统里有单个master和众多chunkserver。Master存储metadata,chunkserver存储实际文件数据。文件被分割成固定大小的chunk (e.g.64MB),每个chunk有一个全局唯一的ID,每个会默认备份3份到不同的机器上,以一个普通的linux文件的形式。

一个读请求的执行流程:客户端发送读请求到master,请求内容包括文件名和offset
--> master返回相应的chunk ID,以及它的若干备份都在哪些机器上
--> 客户端向最近的chunkserver发起读请求,开始传输数据

Master的设计

单master的设计显然带来了single point failure的问题。作者的解释是:1.这样大大减化了设计和实现 2.实际数据直接在客户端和chunkserver间交流,所以单master不会成为bottleneck master内存里的信息主要是每个文件名对应哪些chunk id,以及每台chunkserver上有哪些chunk。只有第一种信息是持久化到本地硬盘以及log中的,第二种信息是持续向chunkserver询问得来,因为chunkserver自己最知道自己的状态,免去了一层没必要的不一致的可能性。

上面提交的master的log非常重要,所以我们要保证先把操作写入log并备份到远程机器后,再对用户可见,和数据库的write-ahead log是一个意思。为了加速挂掉之后的重启速度,会定期作checkpoint,这样重启时只要重跑自上次checkpoint以后的log就行了。

一致性模型(consistency model)

(注:这一块非常重要也比较复杂,以下所述仅为我自己对论文的理解。我大略查了一下网上的资料,细节之处大家的理解并不完全一致。因为没有代码作参考,很难说哪个才是完全对的。)

出于性能考虑,GFS并不是强一致性的,它做出的保证如下:

文件metadata的操作(比如改名)是原子操作,这一点由master里仔细加锁就可以保证。

文件内容的修改有两种,普通写和record append,见下图:(注:“普通写”在论文里简单地用write表示,我在下文称其为“改写”以示和record append的区别)

根据前文所述的假设,record append操作占了绝大多数,改写操作很少。

论文里定义了几种不同的状态,consistent是指多个replica的数据是相同的,但可能内容来自多个并发写入的交叉,不是真正正确的内容,defined是指consistent并且不包含交叉。下文会详述为什么会出现交叉。

为了保证并发写时,各个replica数据的修改顺序是一致的,我们需要指定某个replica为primary,其他的为secondary。在GFS里这个由lease机制来实现。

具体来说,对每个chunk,master负责发放一个lease给其中一个replica,它就变成了primary。lease的有效时间是60秒,但primary可以申请延长,通常情况下都会得到master的批准。当有写请求来时,master会告诉客户端谁是primary,于是客户端总会把写请求发给primary;如果secondary收到了写请求(因为客户端cache了过期信息等原因),它会拒绝执行并通知客户端。master有权提前收回release;如果primary和master无法连通了,lease会在60秒后自动过期。有了lease机制,正常情况下对同一chunk的所有写操作都会发给同一个primary,这样写入的顺序就可以唯一确定下来。

另外,考虑到写入的数据量通常比较大,所以我们把数据流和指令流分开。先把数据发给各个replica,写入本地临时数据区,再传送一个“写入”的指令,这样可以尽可能地避免数据不一致。

我们先来看改写的情况:

  1. 客户端向master询问欲写入chunk的primary是谁,以及所有的secondary是谁。(如果这时没有人有lease,master就会指定一个。)
  2. master返回相关信息,客户端会把信息存入cache,以后就不再重复询问master以节省开销,直到该primary说自己已不再持有lease。
  3. 客户端把数据发给所有replica,replica们会把数据暂存下来,此时并不会真的写入指定位置。
  4. 当所有replica都报告数据收到后,客户端发写入指令给primary。primary会给这个指令定一个序列号(所以当同时收到多个请求时,在primary这里确定顺序)。primary依序列号修改自己的那份数据。
  5. primary把写入指令和序列号发给secondary,secondary都依同样序列号修改自己的数据。
  6. 当primary收到所有人的回复时,返回成功信息给客户端。
  7. 如果有的secondary成功了,有的失败了,primary会返回失败信息给客户端。此时数据就是不一致的。客户端会发起重试。

这里就出现了一个问题,如果欲改写的区域跨越64MB chunk的边界,那么这个写操作就要写多个chunk,而上文的lease机制不能保证对多个chunk的并发写有一致的修改顺序。假设有两个客户端A/B对同一段跨chunk的区域改写不同的内容,就会出现前一个chunk的写入顺序是先A后B,而后一个是先B后A,于是最终结果就是B的前一半和A的后一半留下来。而在GFS看来这完全是成功的写入。这就是上表中的"consistent but undefined"。

Record append的情况有些相似,不同之处在于客户端不会指定offset,每次写入时会针对待写文件的最后一个chunk。如果发现该chunk剩余空间不足以写入,则把当前chunk用空白补齐(padding),然后把数据写入新的chunk。数据的写入是at-least-once,如果写失败了(i.e.只在部分secondary上写成功),则会在新的末尾重新写一次,这样就会导致上次已经写成功的replica上数据出现两次,而上次写失败的replica上会有一段空白。返回给客户端的是成功写入的offset位置。这就是上表里所谓的"interspersed with inconsistent"。客户端程序需要能正确处理这些情况:对于不正确的数据,可以在每个记录开头加一个magic number,或者加一个checksum之类;对于重复的数据,需要客户端判重,比如在记录里加一个unique id。

杂项

GFS没有做数据内容的cache,因为按照预想的应用场景,读操作大多是整个文件从头到尾扫一遍,working set太大,cache能提供的帮助有限。

GFS还提供了snapshot操作,利用copy-on-write机制实现文件或目录的快速复制。

用chunk版本号实现过期数据的检查。

用checksum检验硬盘数据是否被偶然修改(坏道?宇宙射线?)

论文里没有明说,但我觉得master挂掉后是会转为readonly模式,需要人工介入来把backup提升成master。(Google里真正用的应该是更nb的了。)

有一些参考资料介绍了GFS的沿革:

queue.acm.org/detail.cf
highscalability.com/blo

Tuesday, October 12, 2021

RIG stock: Crude oil price | Supply shock

 

Moody’s: Oil Industry Must Spend $542 Billion To Avoid Supply Shock

  1. Annual upstream investment crashed by around 30 percent in 2020 and has only slightly recovered since, according to the credit rating agency
  2. Global annual upstream spending needs to increase by as much as 54 percent to $542 billion if the oil market is to avert the next supply shortage shock

Global annual upstream spending needs to increase by as much as 54 percent to $542 billion if the oil market is to avert the next supply shortage shock, Moody’s said in a report on Thursday carried by Bloomberg.

Currently, oil exploration and production (E&P) companies around the world are underinvesting in supply as they continue to keep capital expenditure (capex) low after the 2020 price crash and crisis, Moody’s notes.

Annual upstream investment crashed by around 30 percent in 2020 and has only slightly recovered since, according to the credit rating agency.

Most producers continue to stick to conservative capital budgets for 2022, but slight growth can be expected as commodity prices jump, Moody’s analysts say.

Yet, “Our analysis demonstrates that upstream companies will need to increase their spending considerably for the medium term to fully replace reserves and avoid declines in future production,” Moody’s Vice President Sajjad Alam said in a statement.

This year, spending is expected at $352 billion, while medium-term annual investment has to grow to $542 billion to keep with the demand returning from the pandemic slump, according to Moody’s.

“The industry will need to spend significantly more, especially if oil and gas demand keeps climbing beyond pre-pandemic levels through 2025,” Moody’s analysts wrote in the report, as carried by Bloomberg.

Underinvestment in upstream projects is a major wild card for oil markets going forward, analysts and industry officials say.

The oil industry is “massively underinvesting” in supply to meet growing demand, which is set to return to pre-COVID levels as soon as the end of 2021 or early 2022, Greg Hill, president of U.S. oil producer Hess Corp, said last week.

Last year, global upstream investment sank to a 15-year low of $350 billion, down from around $600 billion before the pandemic, according to estimates by Wood Mackenzie from earlier this year.

By Tsvetana Paraskova for Oilprice.com

Bigtable: A Distributed Storage System for Structured Data > Building blocks

Oct. 12, 2021

Here are my notes:

  1. GFS - distributed Google File system (GFS) - store log and data files
  2. Bigtable processes often share the same machines with processes from other applications
  3. Bigtable depends on a cluster management system for scheduling jobs, managing resources on shared machines, dealing with machine failures, and monitoring machine status
  4. Google SSTable file format -  internally to store Bigtable data
  5. An SSTable provides a persistent, ordered immutable map from keys to values, where both keys and values are arbitrary byte strings
  6. A block index (stored at the end of the SSTable) is used to locate blocks; the index is loaded into memory when the SSTable is opened.

Bigtable is built on several other pieces of Google infrastructure. Bigtable uses the distributed Google File System (GFS) [17] to store log and data files. A Bigtable cluster typically operates in a shared pool of machines that run a wide variety of other distributed applications, and Bigtable processes often share the same machines with processes from other applications. Bigtable depends on a cluster management system for scheduling jobs, managing resources on shared machines, dealing with machine failures, and monitoring machine status. The Google SSTable file format is used internally to store Bigtable data. An SSTable provides a persistent, ordered immutable map from keys to values, where both keys and values are arbitrary byte strings. Operations are provided to look up the value associated with a specified key, and to iterate over all key/value pairs in a specified key range. Internally, each SSTable contains a sequence of blocks (typically each block is 64KB in size, but this is configurable). A block index (stored at the end of the SSTable) is used to locate blocks; the index is loaded into memory when the SSTable is opened. A lookup can be performed with a single disk seek: we first find the appropriate block by performing a binary search in the in-memory index, and then reading the appropriate block from disk. Optionally, an SSTable can be completely mapped into memory, which allows us to perform lookups and scans without touching disk.

Bigtable relies on a highly-available and persistent distributed lock service called Chubby [8]. A Chubby service consists of five active replicas, one of which is elected to be the master and actively serve requests. The service is live when a majority of the replicas are running and can communicate with each other. Chubby uses the Paxos algorithm [9, 23] to keep its replicas consistent in the face of failure. Chubby provides a namespace that consists of directories and small files. Each directory or file can be used as a lock, and reads and writes to a file are atomic. The Chubby client library provides consistent caching of Chubby files. Each Chubby client maintains a session with a Chubby service. A client’s session expires if it is unable to renew its session lease within the lease expiration time. When a client’s session expires, it loses any locks and open handles. Chubby clients can also register callbacks on Chubby files and directories for notification of changes or session expiration.

Bigtable uses Chubby for a variety of tasks: to ensure that there is at most one active master at any time; to store the bootstrap location of Bigtable data (see Section 5.1); to discover tablet servers and finalize tablet server deaths (see Section 5.2); to store Bigtable schema information (the column family information for each table); and to store access control lists. If Chubby becomes unavailable for an extended period of time, Bigtable becomes unavailable. We recently measured this effect in 14 Bigtable clusters spanning 11 Chubby instances. The average percentage of Bigtable server hours during which some data stored in Bigtable was not available due to Chubby unavailability (caused by either Chubby outages or network issues) was 0.0047%. The percentage for the single cluster that was most affected by Chubby unavailability was 0.0326%.

Read more carefully | SSTable 

A block index (stored at the end of the SSTable) is used to locate blocks; the index is loaded into memory when the SSTable is opened. A lookup can be performed with a single disk seek: we first find the appropriate block by performing a binary search in the in-memory index, and then reading the appropriate block from disk. Optionally, an SSTable can be completely mapped into memory, which allows us to perform lookups and scans without touching disk.

  • A block index (stored at the end of the SSTable) is used to locate blocks; 
  • the index is loaded into memory when the SSTable is opened. 
  • A lookup can be performed with a single disk seek: 
  • we first find the appropriate block by performing a binary search in the in-memory index, and then reading the appropriate block from disk. 
  • Optionally, an SSTable can be completely mapped into memory, which allows us to perform lookups and scans without touching disk.
Google block index - stored at the end of the SSTable  - locate blocks
SSTable - opened, the index is loaded into memory 
a lookup - a single disk seek 
find the appropriate block by performing a binary search in the in-memory index 
reading the appropriate block from disk 
SSTable - mapped into memory, lookups and scans without touching disk

Read more carefully | Bigtable and Chubby 

Bigtable uses Chubby for a variety of tasks: to ensure that there is at most one active master at any time; to store the bootstrap location of Bigtable data (see Section 5.1); to discover tablet servers and finalize tablet server deaths (see Section 5.2); to store Bigtable schema information (the column family information for each table); and to store access control lists. If Chubby becomes unavailable for an extended period of time, Bigtable becomes unavailable. We recently measured this effect in 14 Bigtable clusters spanning 11 Chubby instances. The average percentage of Bigtable server hours during which some data stored in Bigtable was not available due to Chubby unavailability (caused by either Chubby outages or network issues) was 0.0047%. The percentage for the single cluster that was most affected by Chubby unavailability was 0.0326%.

  • Chubby tasks - master slave protocol - one active master at any time 
  • Store the bootstrap location of Bigtable data 
  • discover tablet servers and finalize tablet server deaths 
  • store Bigtable schema information (the column family information for each table)
  • store access control lists
Actionable items

  1. Read GFS paper - GHEMAWAT, S., GOBIOFF, H., AND LEUNG, S.-T. The Google file system. In Proc. of the 19th ACM SOSP (Dec. 2003), pp. 29–43.
  2. Read Chubby paper - BURROWS, M. The Chubby lock service for loosely coupled distributed systems. In Proc. of the 7th OSDI (Nov. 2006).
  3. Read paper: CHANDRA, T., GRIESEMER, R., AND REDSTONE, J. Paxos made live — An engineering perspective. In Proc. of PODC (2007).
  4. Review Chinese article about Paxos paper



BigTable > LSM - Log structured-merge tree

 

作者:Bowen Xiao
链接:https://www.zhihu.com/question/19887265/answer/1695928779
来源:知乎
著作权归作者所有。商业转载请联系作者获得授权,非商业转载请注明出处。

另外吐槽下题目,感觉“算法”二字把LSM渲染成为了数学一样晦涩难懂的程度。更关键它也不是传统意义上的“算法”,顶多算是方法。

和其他索引结构对比下会好理解一点。BTree是第一个比较成功的索引结构。它把数据组织成Page,非叶子节点存放类似指针的值,指向自己的子节点,叶子节点存Value。它比较成功的一点是充分利用磁盘I/O的带宽。

它最本质的问题是系统设计者被传统思路束缚住了手脚:更新数据必须要同时覆盖旧的数据,也就是同一份数据只能存一份。要更新一个值,首先必须要通过查询找到这个值的Page。这带来了无数的问题,比如更新的时候崩溃了,新值旧值参杂在一起,怎么办?(于是有了Write Ahead Log和Double Write,有一些思路也是用Copy-On-Write来解决)比如范围查询怎么做?如果要一个个随机I/O搬回内存可太慢了(于是有了B+ Tree叶子节点的sibiling指针)。又比如某些更新涉及多个Page,很可能发生分裂,如果分裂的时候崩溃了怎么办?(可能有一个孤儿Page,没有被其他任何页所指向)。又比如B Tree是按页为单位,所以即使只有几个Bytes的更改,也要忍受写整个页(4K Bytes)的开销。又比如Page里的内部碎片怎么利用?又比如B Tree的锁都得做在节点中,这让并发控制比较难实现。

相比B/B+ Tree,LSM Tree的设计理念完全不同:更新数据不需要覆盖旧的数据,有多份也不用担心,反正只取最新的值。这一看就能知道肯定写效率会提升:想象一下,你有个纷繁复杂的账本,每次update都直接写在最新的一页而不是翻回前面找到那条记录改掉,多香呀!于是乎实现上可糙多了。它不用太担心崩溃恢复的问题,因为它的新值和旧值本来就物理隔离。感觉它的实现上主要是Trade-Off,方式大家都大差不差,无非就是什么时候做Compaction,Compaction是leveled还是tiered,按行组织的数据什么时候转为列存(或者不转,一般用于AP数据库)等等。LSM最大的问题是后台Compaction线程会和写入线程共同争抢磁盘的I/O资源。

 

Transocean (RIG) set to massively outperform the stock market

Oct. 12, 2021

Here is the link. 

Transocean (RIG) set to massively outperform the stock market due to the following variables: 1. Moody's forecast a significant shortfall in CAPEX for upstream operations and foresees the need for CAPEX to increase by 54% directly impacting companies such as Transocean. 2. Offshore utilization is anticipated to reach 80% by January 2022 meaning there is a significant shortfall in available rigs. 3. The push for a ban on new drilling on federal land will likely cause larger integrate oil companies to rethink their strategy likely leading to more spending pushed into international waters. 3. US government policy towards energy is hawkish. 4. Transocean continue to reinvest and have the newest fleet on the market. I hope you enjoy my stock analysis of Transocean, ticker RIG, thank you as always for your support and have a wonderful day!

BigTable > Log-structure merge-tree (LSM tree) > LSM原理解读 | 20 minutes to warmup

 LSM原理解读  本文来源:码农网 本文链接:https://www.codercto.com/a/49359.html 

内容简介:06年,Google发表了BigTable论文,从此推开了大数据时代的大门。为什么又提及BigTable,是因为这篇论文中使用了LSM这种数据结构。LSM: Log Structured-Merge Tree(日志结构合并树),是一种先于BigTable出现的文件组织方式,最早可以追溯到1996年 Patrick O'Neil等人的论文。在BigTable出现之后开始被重视被广泛应用,比较有名的产品就是Hbase、Cassandra、Leveldb等。

06年,Google发表了BigTable论文,从此推开了大数据时代的大门。

为什么又提及BigTable,是因为这篇论文中使用了LSM这种数据结构。

LSM: Log Structured-Merge Tree(日志结构合并树),是一种先于BigTable出现的文件组织方式,最早可以追溯到1996年 Patrick O'Neil等人的论文。在BigTable出现之后开始被重视被广泛应用,比较有名的产品就是Hbase、Cassandra、Leveldb等。

  • 用一句话先简单概括其原理

把一颗大树拆分成N棵小树,数据先写入内存中,随着小树越来越大,内存的小树会flush到磁盘中。 磁盘中的树定期做合并操作,合并成一棵大树,以优化读性能。

  • 先用个简单的结构看一下

 最简单的LSM tree是两层树状结构C0,C1。 C0比较小,驻留在内存,当C0超过一定的大小, 一些连续的片段会从C0移动到磁盘中的C1中,这是一次merge的过程。在实际的应用中, 一般会分为更多的层级(level), 而层级C0都会驻留在内存中。

背景

LSM被设计来提供比传统的B+树或者ISAM(Indexed Sequential Access Method)更好的写操作吞吐量,通过消去随机的本地更新操作来达到这个目标。之所以说这是一个好方法,是因为对于存储来讲问题的本质还是磁盘随机操作慢,顺序读写快。顺序读写磁盘快于随机读写主存,而且快至少三个数量级。

如果对写操作的吞吐量敏感,一个好的办法是简单的将数据添加到文件。这个策略经常被使用在日志或者堆文件,因为他们是完全顺序的,所以可以提供非常好的写操作性能。

同时他们也有一些缺点,从日志文件中读一些数据将会比写操作需要更多的时间,需要倒序扫描,直接找到所需的内容。

这说明日志仅仅适用于一些简单的场景:1. 数据是被整体访问,像大部分 数据库 的WAL(write-ahead log) 2. 知道明确的offset,比如在Kafka。

与此同时我们也有多种方法来支持复杂的读场景,比如二分查找、哈希、B+数等。这些方式虽提升了读的性能,却丢失了日志文件很好的写性能。

不管怎样,随机的操作要尽量减少。所以LSM就出现了,保持了日志文件写性能,以及微小的读操作性能损失。本质上就是让所有的操作顺序化。

算法

LSM将之前使用一个大的查找结构,变换为将写操作顺序的保存到一些相似的有序文件(Sorted Strings Table,简称sstable)中。所以每个文件包 含短时间内的一些改动。因为文件是有序的,所以之后查找也会很快。文件是不可修改的,他们永远不会被更新,新的更新操作只会写到新的文件中。通过周期性的合并这些文件来减少文件个数。 此时可以回到上面看图 。

更具体的看,当一些更新操作到达时,他们会被写到内存缓存(也就是memtable)中,memtable使用树结构来保持key的有序,在大部分的实现中,memtable会通过写WAL的方式备份到磁盘,用来恢复数据,防止数据丢失。当memtable数据达到一定规模时会被刷新到磁盘上的一个新文件,重要的是系统只做了顺序磁盘读写,因为没有文件被编辑,新的内容或者修改只用简单的生成新的文件。

因为比较旧的文件不会被更新,重复的纪录只会通过创建新的纪录来覆盖,这也就产生了一些冗余的数据。所以系统会周期的执行合并操作(compaction)。 合并操作选择一些文件,并把他们合并到一起,移除重复的更新或者删除纪录,同时也会删除上述的冗余。更重要的是,通过减少文件个数的增长,保证读操作的性能。因为sstable文件都是有序结构的,所以合并操作也是非常高效的。

当一个读操作请求时,系统首先检查内存数据(memtable),如果没有找到这个key,就会逆序的一个一个检查sstable文件,直到key 被找到。因为每个sstable都是有序的,所以查找比较高效(O(logN)),但是读操作会变的越来越慢随着sstable的个数增加,因为每一个 sstable都要被检查。

所以,读操作比其它本地更新的结构慢,幸运的是,有一些技巧可以提高性能。最基本的的方法就是页缓存(也就是leveldb的 TableCache,将sstable按照LRU缓存在内存中)在内存中,减少二分查找的消耗。

即使有每个文件的索引,随着文件个数增多,读操作仍然很慢。通过周期的合并文件,来保持文件的个数,因些读操作的性能在可接收的范围内。即便有了合 并操作,读操作仍然会访问大量的文件,大部分的实现通过布隆过滤器来避免大量的读文件操作, cassandra就采用了Bloom filter的方式。

思考

我们看到LSM有更好的写性能,同时LSM还有其它一些好处。sstable文件是不可修改的,这让对他们的锁操作非常简单。一般来说,唯一的竞争资源就是 memtable,相对来说需要相对复杂的锁机制来管理在不同的级别。

而如scylladb,通过异步IO,非阻塞线程,线程和线程之间通过消息通讯,没有内存的共享,没有互斥锁,也没有原子操作来大幅的提高了cpu的效率。

2016年,Lanyue Lu等人发布了WiscKey论文,这篇论文提出了一种新的设计 WiscKey: Separating Keys from Values in SSD-conscious Storage,专门为SSD所优化,将key和value分别存储以减少I/O放大。key存在LSM tree中, value存在WAL中,叫做value log。

通常情况下,key比较小,所以LSM tree比较小,当获取value值的时候,再从SSD存储中读取。



本文来源:码农网
本文链接:https://www.codercto.com/a/49359.html

RIG stock: Inside purchase

Oct. 12, 2021

There were 15 insider buys by 3 insiders in June for 25.8 million shares at prices ranging from $3.97 to $4.40.

Jeremy Thigpen. “These better than anticipated results were driven largely by our continued focus on operational excellence, as evidenced by our strong uptime performance, which resulted in revenue efficiency of 98 percent.”

RIG is heavily manipulated by Short Term Traders or MM. I listened to the earnings call, there was nothing in the earnings for RIG to be down 10%. Longs just wait, RIG will be green soon. There are too many quick profit seekers out there. They eventually loose money.

RIG has $1B EBITDA and $ 2B Market Cap. RIG shares are way underpriced. This is Wall Street folks.

RIG has $1B EBITDA and $ 2B Market Cap. RIG shares are way underpriced. This is Wall Street folks.

Director https://www.nasdaq.com/market-activity/insiders/mohn-frederik-wilhelm-104430

bought about 100million bucks worth of stock in OPEN market at average 4.20 in June 2021.

He now owns about 15% of RIG

Keep it simple. Going forward RIG will trade with the whole energy sector. Last 3 weeks the whole sector got killed. Once delta is under control next 3 weeks the whole sector will fly. In my opinion earnings will not make much of a difference going forward. Look at XOM, MTDR and CVX. They had solid quarters yet they went down. What matters is where oil is heading next 6 months. I say $80





大数据那些事(2):三驾马车 | GFS

Oct. 12, 2021

Here is the article. 

作者:飞总
链接:https://zhuanlan.zhihu.com/p/24382357
来源:知乎
著作权归作者所有。商业转载请联系作者获得授权,非商业转载请注明出处。

微信公众号:飞总的IT世界面面观。

但凡是要开始讲大数据的,都绕不开最初的Google三驾马车:Google File System(GFS), MapReduce,BigTable。如果我们拉长时间轴到20年为一个周期来看呢,这三驾马车到今天的影响力其实已然不同。MapReduce作为一个有很多优点又有很多缺点的东西来说,很大程度上影响力已经释微了。BigTable以及以此为代表的各种KeyValue Store还有着它的市场,但是在Google内部Spanner作为下一代的产品,也在很大程度上开始取代各种各样的的BigTable的应用。而作为这一切的基础的Google File System,不但没有任何倒台的迹象,还在不断的演化,事实上支撑着Google这个庞大的互联网公司的一切计算。



作为开源最为成功的Hadoop Ecosystem来说,这么多年来风起云涌,很多东西都在变。但是HDFS的地位却始终很牢固。尽管各大云计算厂商都基于了自己的存储系统实现了HDFS的接口,但是这个文件系统的基本理念至今并无太多改变。

那么GFS到底是什么呢?简单一点来说它是一个文件系统,类似我们的NTFS或者EXT3,只是这个文件系统是建立在一堆的计算机的集群之上,用来存储海量数据的。一般来说,对文件系统的操作无非是读或者写,而写通常又被细分成update和append。Update是对已有数据的更新,而append则是在文件的末尾增加新的数据。Google File System的论文发表于2003年,那个时候Google已经可以让这个文件系统稳定的运行在几千台机器上,那么今天在我看来稳定的运行在几万台机器上并非是什么问题。这是非常了不起的成就,也是Hadoop的文件系统至今无非达到的高度。

GFS的设计理念上做了两个非常重要的假设,其一是这个文件系统只处理大文件,一般来说要以TB或者PB作为级别去处理。其二是这个文件系统不支持update只支持append。在这两个假设的基础上,文件系统进一步假设可以把大文件切成若干个chunk,本文上面的图大致上给了GFS的一个基本体系框架的解释。

Chunk server是GFS的主体,它们存在的目的是为了保存各种各样的chunk。这些chunk代表了不同文件的不同部分。为了保证文件的完整性不受到机器下岗的影响,通常来说这些chunk都有冗余,这个冗余标准的来说是3份。有过各种分析证明这个三份是多门的安全。在我曾经工作的微软的文件系统实现里面也用过三份这样的冗余。结果有一次data center的电源出问题,导致一个集装箱的机器整个的失联(微软的数据中心采用了非常独特的集装箱设计)。有一些文件的两个或者三个拷贝都在那个集装箱对应的机器上,可以想象,这也同样导致了整个系统的不可用。所以对于这三个拷贝要放哪里怎么去放其实是GFS里比较有意思的一个问题。我相信Google一定是经过了很多的研究。

除了保存实际数据的chunk server以外,我们还需要metadata manager,在GFS里面这个东西就是master了。按照最初的论文来说,master是一个GFS里面唯一的。当然后续有些资料里有提到GFS V2的相关信息表明这个single point bottleneck 在Google的系统演进中得到了解决。Master说白了就是记录了各个文件的不同chunk以及它们的各自拷贝在不同chunk server上的区别。

Master的重要性不言而喻。没有了metadata的文件系统就是一团乱麻。Google的实现实际上用了一个Paxos协议,倘若我的理解是正确的话。Paxos是Lamport提出来的用来解决在不稳定网络里面的consensus的一个协议。协议本身并不难懂,但是论文的证明需要有些耐心。当然最重要的,我自己从来没有实现过这个协议。但是就我能看到的周围实现过的人来说,这个东西很坑爹。Paxos干的事情是在2N+1台机器里保持N的冗余。简单一点说,挂掉N台机器这个metadata service依然可以使用。协议会选出一个master服务,而其他的则作为shadow server存在。一旦master挂了大家会重新投票。这个协议很好的解决了High Availability的问题。很不幸的是,抄袭的Hadoop 文件系统使用的是一个完全不同的方案。这个我们在讲到Hadoop的时候再说。

对GFS的访问通过client,读的操作里,client会从master那边拿回相应的chunk server,数据的传输则通过chunk server和client之间进行。不会因此影响了master的性能。而写的操作则需要确保所有的primary以及secondary都写完以后才返回client。如果写失败,则会有一系列的retry,实在不行则这些chunk会被标注成garbage,然后被garbage collection。

Master和chunk server之间会保持通信,master保持着每个chunk的三个copy的实际位置。当有的机器掉线之后,master如果有必要也会在其他的机器上触发额外的copy活动以确保冗余,保证文件系统的安全。

GFS的设计非常的值得学习。系统支持的目标文件以及文件的操作非常的明确而简单。对待大规模不稳定的PC机构成的data center上怎么样建立一个稳定的系统,对data使用了复制,而对metadata则用了Paxos这样的算法的实现。这个文件系统体现了相当高水准的系统设计里对方方面面trade off的选择。有些东西只有自己做过或者就近看人做过才能真正感受到这系统设计的博大精深。故而对我个人而言,我对GFS的论文一直是非常的推崇,我觉得这篇论文值得每个做系统的人反复的读。

大数据的那些事(3):三驾马车 | MapReduce

Oct. 12, 2021

Here is the article. 

作者:飞总
链接:https://zhuanlan.zhihu.com/p/24382810
来源:知乎
著作权归作者所有。商业转载请联系作者获得授权,非商业转载请注明出处。

微信公众号:飞总的IT世界面面观。

在Google的三驾马车里面,Google File System是永垂不朽的,也是基本上没有人去做什么进一步的研究的。BigTable是看不懂的,读起来需要很多时间精力。唯独MapReduce,是霓虹灯前面闪烁的星星,撕逼战斗的主角,众人追捧和喊打的对象。自从MapReduce这个词出来以后,不知道有多少篇论文发表出来,又不知道有多少口诛笔伐的文章。

作为论文来说MapReduce严格的来讲不能算作一篇论文,因为它讲述了两件不同的事情。其一是一个叫做MapReduce的编程模型。其二是大规模数据处理的体系架构的实现。这篇论文将两者以某种方式混杂在一起来达到不可告人的目的,并且把这个体系吹得非常的牛,但是却并没有讨论一些Google内部造就知道的局限性,以我对某狗的某些表现来看,恐怕我的小人之心觉得有意为之的可能性比较大。

因此当智商比较低的Yahoo活雷锋抄袭MapReduce的时候弄出的Hadoop是不伦不类,这才有了后来Hadoop V2以及Yarn的引进。当然这是后话。作为同样抄袭对象的微软就显得老道很多。微软内部支撑大数据分析的平台Cosmos是狠狠的抄袭了Google的File system却很大程度上摒弃了MapReduce这个框架。当然这些都是后话,只能留待以后的篇幅慢慢展开。

我们先看看作为编程模型的MapReduce。所谓MapReduce的意思是任何的事情只要都严格遵循Map Shuffle Reduce三个阶段就好。其中Shuffle是系统自己提供的而Map和Reduce则用户需要写代码。Map是一个per record的操作。任何两个record之间都相互独立。Reduce是个per key的操作,相同key的所有record都在一起被同时操作,不同的key在不同的group下面,可以独立运行。这就像是说我们有一把大砍刀,一个锤子。世界上的万事万物都可以先砍几刀再锤几下,就能搞定。至于刀怎么砍,锤子怎么锤,那就算个人的手艺了。

从计算模型的角度来看,这个模型极其的粗糙。所以现在连Google自己都不好意思继续鼓吹MapReduce了。从做数据库的人的角度来看这无非是一个select一个groupby,这些花样197x的时候在SystemR里都被玩过了。数据库领域玩这些花样无数遍。真看不出有任何值得鼓吹的道理。因此,在计算模型的角度上来说,我觉得Google在很大程度上误导和夸大了MapReduce的实际适用范围,也可能是自己把自己也给忽悠了。

在Google内部MapReduce最大的应用是作为inverted index的build的平台。所谓inverted index是information retrieval里面一个重要的概念,简单的讲是从单词到包含单词的文本的一个索引。我们搜索internet,google需要爬虫把网页爬下来,然后建立出网页里面的单词到这个网页的索引。这样我们输入关键字搜索的时候,对应的页面才能出来。也正因为是这样,所以Google的论文里面用了word count这个例子。下图是word count的MapReduce的一个示意图。



然而我们需要知道的是,Google后来公布的信息显示它的广告系统是一直运行在MySQL的cluster的,该做join的时候也是做join的。MapReduce作为一个编程模型来说,显然不是万能的药。可是因为编程模型涉及的是世界观方法论的问题。于是催生了无数篇论文,大致的套路都是我们怎么样用MapReduce去解决这个那个问题。这些论文催生了无数PhD,帮助很多老师申请到了很多的钱。我觉得很大程度上都掉进了google的神话和这个编程模型的坑。

MapReduce这篇论文的另外一个方面是系统实现。我们可以把题目写成:如何用一堆廉价PC去稳定的实现超大规模的并行数据处理。我想这无疑可以体现出这篇论文真正有意义的地方。的确,数据库的工业界和学术界都玩了几十年了,有哪个不是用高端的机器。在MapReduce论文出来的那个时候,谁能处理1个PB的数据我给谁跪了。但是Google就能啊。我得意的笑我得意的笑。所以Google以它十分牛逼的数据处理平台,去吹嘘那个没有什么价值的编程模型。而数据库的人以攻击Google十分不行的编程模型,却故意不去看Google那个十分强悍的数据处理平台。这场冯京对马凉的比赛,我觉得毫无意义。

那么我们来看看为什么Google可以做到那么大规模的数据处理。首先这个系统的第一条,很简单,所有的中间结果可以写入到一个稳定的,不因为单机的失败而不能工作的分布式海量文件系统。GFS的伟大可见一斑。没有GFS,玩你妹的MapReduce。没有一个database厂商做出过伟大的GFS,当然也就没办法做出这么牛叉的MapReduce了。

这个系统的第二条也很简单,能够对单个worker进行自动监视和retry。这一点就使得单个节点的失败不是问题,系统可以自动的进行管理。加上Google一直保持着绝不泄密的资源管理系统Borg。使得Google对于worker能够进行有效的管理。

Borg这个系统存在有10多年了,但是Google故意什么都不告诉大家,论文里也假装没有。我第一次听说是几个从Google出来的人在Twitter想重新搞这样一个东西。然而一直到以docker为代表的容器技术的出现,才使得大家知道google的Borg作为一个资源管理和虚拟化系统到底是怎么样做的。而以docker为代表的容器技术的出现也使得Borg的优势不存在了。所以Google姗姗来迟的2015年终于发了篇论文。我想这也是Yahoo这个活雷锋没有抄好,而HadoopV2必须引入Yarn的很重要的原因。

解释这么多,其实是想说明几点,MapReduce作为编程模型,是一个很傻的模型。完全基于MapReduce的很多project都不太成功。而这个计算模型最重要的是做inverted index build,这就使得Google长久以来宣扬的Join没意义的论调显得很作。另外随着F1的披露,大家知道Google的Ads系统实际上长期运行在MySQL上,这也从侧面反应了Google内部的一些情况和当初论文的高调宣扬之间的矛盾。

Google真正值得大家学习的是它怎么样实现了大规模数据并发的处理。这个东西说穿了,一是依赖于一个很牛的文件系统,二是有着很好的自动监控和重试机制。而MapReduce这个编程模型又使得这两者的实现都简化了。然而其中很重要的资源管理系统Borg又在当初的论文里被彻底隐藏起来了。我想,随着各种信息的披露,我只能说一句,你妹的。

MapReduce给学术界掀起了一片灌水高潮,学术界自娱自乐的精神实在很值得敬佩。

BigTable > Log structured-merge tree (LSM) tree | 大数据, 复活的LSM-Tree--BigTable的系统实现

Here is the article.  

我的微信公众号“飞总的IT世界面面观”上有后续的大数据那些事,欢迎关注。

BigTable是一个非常复杂的系统,发表的论文面面俱到,但是每个方面都写得并不是很清楚。所幸Google开源了LevelDB这个Key-Value Store。这个项目的作者是Jeff Dean和Sanjay Ghemawat,被认为很大程度上重复使用了BigTable在单个节点上的实现。LevelDB为我们对BigTable的实现提供了重要的学习资料。

在BigTable的实现上,一个BigTable的cluster由一个client library,一个Master server和很多个的Tablet Server组成。按照论文的说法,一个大的BigTable会被分成若干个大小大致在100MB到200MB的tablets,而这些tablets 会被分配到一些Tablet Server上去给client 提供服务来。Tablet Server对超过200MB的tablet灰进行切分,对小于100MB的则会进行合并。

系统运行过程中的Tablet server的数量不是固定的,可以根据实际上的工作负载来增加或者减少,这方面的工作是Master server来控制的。Tablet server并不存储实际的文件,而是作为一种service和proxy来访问存在Google File System里的实际的tablet们。

和Tablet server不一样的是,Master Server始终都存在。Master server存在主要是把tablet分配给Tablet server,增加或者减少Tablet Server,并且负责去平衡不同的Tablet server的load。如果用户要创建一个新的Table,或者对已经有的Table做改动的话,譬如增加新的column family等,都是通过Master server来完成的。Master server通过Tablet server来实现对tablet的间接操作,本身并不负责对任何Tablet 的管理。

和大家直观上想象的不同,当一个client要访问BigTable的时候,client并不需要和Master server进行交流。这就保证了Master server的load并不是很重。那么,client是怎么样实现对BigTable的访问的呢? 这是BigTable比较精密的difference。这需要用到Chubby。

Chubby是一个highly distributed lock service。可以认为开源Zookeeper是一个Chubby的copycat。但是虽然说是CopyCat,实际上Chubby实现的是一个Paxos协议,而Zookeeper实现的是它的变种Zab。这方面我们不展开细讲。 具体的情况请阅读 The Chubby lock service for loosely-coupled distributed systems。

Chubby实现的是一个类似文件系统那样的目录结构。使用者可以访问这些文件来获得对被访问对象的锁。按照BigTable论文的说法,Chubby的用处有很多处,包括对Tablet的定位,对Tablet server的监控等等。

对我们来说最重要的是了解client怎么样对数据进行操作。这个操作首先要对Metadata进行访问。这个操作大致上是要通过访问一个三层的结构,其中第一层是一个Chubby file。通过Chubby file可以找到一个叫做Root Tablet的位置,Root Tablet下一层则是所有的metadata tablets。有一点非常的特别,Root Tablet从来不会被partition。所以一般来说client需要通过访问一个Chubby文件,一个Root Tablet和一个Metadata Tablet的三层结构去定位到对应的data。Client 会cache住它找到的metadata,但是cache并非必须。只要Chubby service是活着的,通过这个三层访问总是可以定位到对应的data的tablet的。通过Chubby,Master不需要承担任何datat活着metadata的发现。故而Master server的负担很轻。

论文没有谈及metadata的格式到底是什么,最详细描述基于下面的论文原话: The Meadata Table store ... location of... tablet under a row key that is encoding of the tablet's table identifier and its end row. 但是其实不重要,聪明的engineer总是可以设计出自己的metadata的,比如说HBase就实现了自己的一套metadata。我想这个实现和BigTable应该很不一样。

在BigTable里, SSTable (Sorted Strings Table) 是一个基本的单元。每个 Tablet 有若干个SSTable。论文里面并没有提到SSTable是怎么样实现的。但是根据对开源的LevelDB的学习,我们可以大致上可以描述一下SSTable是怎么实现的。这个实现复活了1993年由UMass Boston退休教授 Patrick O'Neil实现的LSM-Tree。我们把相关的论文放在下面供大家参考:

Patrick O'Neil, Edward Cheng, Dieter Gawlick, and Elizabeth (Betty) O'Neil, "The Log-Structured Merge-Tree," patent granted to Digital Equipment Corporation, December 1993; appeared in ActaInformatica 33, pp. 351-385, June 1996.

看起来这个tree是申请了专利的,专利授权给了已经挂掉的DEC公司。不知道Google用这个Tree是不是侵权,也不知道目前用了LevelDB的,以及Facebook克隆LevelDB做出来的RocketDB的各大公司们有没有获得专利公司的授权。还是说因为年代太久远而DEC公司已经挂掉了无数年了,大家可以肆无忌惮的用了呢?

整个SSTable的实现分为memTable和磁盘上的SSTable。在内存里使用的是skip-list。所以写的操作只是写内存,非常的快。而内存写满之后就会把这个memTable变成一个immutable的memTable。同时开一个新的可以写的memTable另外一个线程则会把这个immutable的内存表变成一个磁盘上的SSTable。当这个转变完成以后immutable的内存表被释放。如此往复磁盘上会产生很多的SSTable。这就需要compact。SSTable是有level1, level2,。。。的。其中进到level1的compact叫做minor compact。后面还有major compact。从level2开始以后任意两棵树的key之间不会有overlap,但是在level1这并不guarantee。

所以我们的一个读操作要读memTable,immutable memTable,level1的tree,和level2以及以下的level的1棵树。这说明读的操作相对写的操作会更贵一些。大家需要注意的是,如果我们访问的是最新的版本,那么有可能会在内存里,所以这个设计对于读的操作主要是优化了新版本的读。对于cool data的读则要慢很多。顺便补充一句,Facebook的copycat RocksDB和LevelDB最主要的区别据说是引入了一个叫做universal compact的东西。当然我没有研究过这个codebase,不清楚universal compact到底有多牛。

当然,就像任何一个类似的系统一样,BigTable的recovery基于log,所有的写操作进内存之前写进log。LevelDB的log format并不是太难懂,是经典的append only的操作。基于log的读写恢复是任何一个系统的基础了。我就不再展开叙述。

说起来非常的有意思,这个LSM-Tree是一个挺了不起的发明,但是波澜不兴的在Industry很多年。发明者Patrick O'Neil在DB界固然是混一个教授到退休,但是也没有获得和他发明相匹配的知名度。也不知道Google的员工们是怎么样从一本并不怎么知名的杂志里把这个东西给挖出来,复活了。有时候我在想,倘若当年没有成名可以理解,现在已经成为了BigTable的基础了,那么作为database research community是不是应该recognize一下。实际上当然是没有。出乎我意料之外的是今年2016的VLDB我居然在印度遇到了Patrick的PhD,实在无法想象这么老了还在带学生。Database community有一些时候对某些人的recognition可能确实是差了一些。比如说Peter Chen,这位在MIT没有拿到tenure在路易斯安那州立大学安家的人,发明了E-R model,为database的发扬光大做出了巨大的贡献。但是我觉得整个community对他的recognition非常的有限。比如说Goetz Graefe,这个在query optimization和query execution做出过巨大贡献的人,发明了Volcano optimizer和经典的Volcano execution model,包括今天大行其道的Exchange operator,居然到今天也没有获得 Sigmod Edgar Codd Innovations Award,我实在是不能理解。

编辑于 2017-01-04