Wednesday, October 6, 2021

BigTable: Cloud Bigtable > Documentation > Guides > Understanding Bigtable performance

 Understanding Bigtable performance

This page describes the approximate performance that Cloud Bigtable can provide under optimal conditions, factors that can affect performance, and tips for testing and troubleshooting Bigtable performance issues.

Bigtable delivers highly predictable performance that is linearly scalable. When you avoid the causes of slower performance described below, each Bigtable node can provide the following approximate throughput, depending on which type of storage the cluster uses:

Storage TypeReads Writes Scans
SSDup to 10,000 rows per secondorup to 10,000 rows per secondorup to 220 MB/s
HDDup to 500 rows per secondorup to 10,000 rows per secondorup to 180 MB/s

These estimates assume that each row contains 1 KB of data.

In general, a cluster's performance scales linearly as you add nodes to the cluster. For example, if you create an SSD cluster with 10 nodes, the cluster can support up to 100,000 rows per second for a typical read-only or write-only workload.

When planning your Bigtable clusters, it is important to think about the trade-off between throughput and latency. Bigtable is used in a broad spectrum of applications, and different use cases can have different optimization goals. For example, for a batch data processing job, you might care more about throughput but less about latency. On the other side, an online service that serves user requests might prioritize lower latency over throughput. As a result, it is important to plan the capacity accordingly.

The numbers in the Performance for typical workloads section are achievable when you prioritize throughput, but the tail latency for Bigtable under such a load might be too high for latency-sensitive applications. In general, Bigtable offers optimal latency when the CPU load for a cluster is under 70%. For latency-sensitive applications, however, we recommend that you plan at least 2x capacity for your application's max Bigtable queries per second (QPS). This capacity ensures that your Bigtable cluster runs at less than 50% CPU load, so it can offer low latency to front-end services. This capacity also provides a buffer for traffic spikes or key-access hotspots, which can cause imbalanced traffic among nodes in the cluster.

Another consideration in capacity planning is storage. The storage capacity of a cluster is determined by the storage type and the number of nodes in the cluster. When the amount of data stored in a cluster increases, Bigtable optimizes the storage by distributing the amount of data across all the nodes in the cluster.

You can determine the storage usage per node by dividing the cluster's storage utilization (bytes) by the number of nodes in the cluster. For example, consider a cluster that has three HDD nodes and 9 TB of data. Each node stores about 3 TB, which is 18.75% of the HDD storage per node limit of 16 TB.

When storage utilization increases, workloads can experience an increase in query processing latency even if the cluster has enough nodes to meet overall CPU needs. This is because the higher the storage per node, the more background work such as indexing is required. The increase in background work to handle more storage can result in higher latency and lower throughput.

For latency-sensitive applications we recommend that you keep storage utilization per node below 60%. If your dataset grows, add more nodes to maintain low latency.

For applications that are not latency-sensitive, you can store more than 70% of the limit, as explained in Storage per node.

Always run your own typical workloads against a Bigtable cluster when doing capacity planning, so you can figure out the best resource allocation for your applications.

Google's PerfKit Benchmarker uses YCSB to benchmark cloud services. You can follow the PerfKitBenchmarker tutorial for Bigtable to create tests for your own workloads. When doing so, you should tune the parameters in the benchmarking config yaml files to make sure that the generated benchmark reflects the following characteristics in your production:

Refer to Testing performance with Bigtable for more best practices.

There are several factors that can cause Bigtable to perform more slowly than the estimates shown above:

  • You read a large number of non-contiguous row keys or row ranges in a single read request. Bigtable scans the table and reads the requested rows sequentially. This lack of parallelism affects the overall latency, and any reads that hit a hot node can increase the tail latency. See Reads and performance for details.
  • The table's schema is not designed correctly. To get good performance from Bigtable, it's essential to design a schema that makes it possible to distribute reads and writes evenly across each table. See Designing Your Schema for more information.
  • The rows in your Bigtable table contain large amounts of data. The performance estimates shown above assume that each row contains 1 KB of data. You can read and write larger amounts of data per row, but increasing the amount of data per row will also reduce the number of rows per second.
  • The rows in your Bigtable table contain a very large number of cells. It takes time for Bigtable to process each cell in a row. Also, each cell adds some overhead to the amount of data that's stored in your table and sent over the network. For example, if you're storing 1 KB (1,024 bytes) of data, it's much more space-efficient to store that data in a single cell, rather than spreading the data across 1,024 cells that each contain 1 byte. If you split your data across more cells than necessary, you might not get the best possible performance. If rows contain a large number of cells because columns contain multiple timestamped versions of data, consider keeping only the most recent value. Another option for a table that already exists is to send a deletion for all previous versions with each rewrite.
  • The cluster doesn't have enough nodes. A cluster's nodes provide compute for the cluster to handle incoming reads and writes, keep track of storage, and perform maintenance tasks such as compaction. You need to make sure that your cluster has enough nodes to satisfy the recommended limits for both compute and storage. Use the monitoring tools to check whether the cluster is overloaded.

    • Compute - If the CPU of your Bigtable cluster is overloaded, adding more nodes can improve performance by spreading the workload across more nodes.
    • Storage - If your storage usage per node has become higher than recommended, you need to add more nodes to maintain optimal latency and throughput, even if the cluster has enough CPU to process requests. This is because increasing storage per node increases the amount of background maintenance work per node. For details, see Trade-offs between storage usage and performance.
  • The Bigtable cluster was scaled up or scaled down recently. After you increase the number of nodes in a cluster to scale up, it can take up to 20 minutes under load before you see a significant improvement in the cluster's performance. When you decrease the number of nodes in a cluster to scale down, try not to reduce the cluster size by more than 10% in a 10-minute period to minimize latency spikes.

  • The Bigtable cluster uses HDD disks. In most cases, your cluster should use SSD disks, which have significantly better performance than HDD disks. See Choosing between SSD and HDD storage for details.

  • There are issues with the network connection. Network issues can reduce throughput and cause reads and writes to take longer than usual. In particular, you might see issues if your clients are not running in the same zone as your Bigtable cluster, or if your clients run outside of Google Cloud.

Because different workloads can cause performance to vary, you should perform tests with your own workloads to obtain the most accurate benchmarks.

Enabling replication will affect the performance of a Bigtable instance. The effect is positive for some metrics and negative for others. You should understand potential impacts on performance before deciding to enable replication.

Replication can improve read throughput, especially when you use multi-cluster routing. Additionally, replication can reduce read latency by placing your Bigtable data geographically closer to your application's users.

Although replication can improve availability and read performance, it does not increase write throughput. A write to one cluster must be replicated to all other clusters in the instance. As a result, each cluster is expending CPU resources to pull changes from the other clusters. Write throughput might actually go down because replication requires each cluster to do additional work.

For example, suppose you have a single-cluster instance, and the cluster has 3 nodes:

Single-cluster instance that has 3 nodes

If you add nodes to the cluster, the effect on write throughput is different than if you enable replication by adding a second 3-node cluster to the instance.

Adding nodes to the original cluster: You can add 3 nodes to the cluster, for a total of 6 nodes. The write throughput for the instance doubles, but the instance's data is available in only one zone:

Single-cluster instance that has 6 nodes

With replication: Alternatively, you can add a second cluster with 3 nodes, for a total of 6 nodes. The instance now writes each piece of data twice: when the write is first received and again when it is replicated to the other cluster. The write throughput does not increase, and might go down, but you benefit from having your data available in two different zones:

Two-cluster instance that has 6 nodes

In these examples, the single-cluster instance can handle twice the write throughput that the replicated instance can handle, even though each instance's clusters have a total of 6 nodes.

When you use multi-cluster routing, replication for Bigtable is eventually consistent. As a general rule, it takes longer to replicate data across a greater distance. Replicated clusters in different regions will typically have higher replication latency than replicated clusters in the same region.

Depending on your use case, you will use one or more app profiles to route your Bigtable traffic. Each app profile uses either multi-cluster or single-cluster routing. The choice of routing can affect performance.

Multi-cluster routing can minimize latency. An app profile with multi-cluster routing automatically routes requests to the closest cluster in an instance from the perspective of the application, and the writes are then replicated to the other clusters in the instance. This automatic choice of the shortest distance results in the lowest possible latency.

An app profile that uses single-cluster routing can be optimal for certain use cases, like separating workloads or to have read-after-write semantics on a single cluster, but it will not reduce latency in the way multi-cluster routing does.

To understand how to configure your app profiles for these and other use cases, see Examples of Replication Settings.

To store the underlying data for each of your tables, Bigtable shards the data into multiple tablets, which can be moved between nodes in your Bigtable cluster. This storage method enables Bigtable to use two different strategies for optimizing your data over time:

  1. Bigtable tries to store roughly the same amount of data on each Bigtable node.
  2. Bigtable tries to distribute reads and writes equally across all Bigtable nodes.

Sometimes these strategies conflict with one another. For example, if one tablet's rows are read extremely frequently, Bigtable might store that tablet on its own node, even though this causes some nodes to store more data than others.

As part of this process, Bigtable might also split a tablet into two or more smaller tablets, either to reduce a tablet's size or to isolate hot rows within an existing tablet.

The following sections explain each of these strategies in more detail.

As you write data to a Bigtable table, Bigtable shards the table's data into tablets. Each tablet contains a contiguous range of rows within the table.

If you have written only a small amount of data to the table, Bigtable will store all of the tablets on a single node within your cluster:

A cluster with four tablets on a single node.

As more tablets accumulate, Bigtable will move some of them to other nodes in the cluster so that the amount of data is balanced more evenly across the cluster:

Additional tablets are distributed across multiple nodes.

If you've designed your schema correctly, then reads and writes should be distributed fairly evenly across your entire table. However, there are some cases where you can't avoid accessing certain rows more frequently than others. Bigtable helps you deal with these cases by taking reads and writes into account when it balances tablets across nodes.

For example, suppose that 25% of reads are going to a small number of tablets within a cluster, and reads are spread evenly across all other tablets:

Out of 48 tablets, 25% of reads are going to 3 tablets.

Bigtable will redistribute the existing tablets so that reads are spread as evenly as possible across the entire cluster:

The three hot tablets are isolated on their own node.

If you're running a performance test for an application that depends on Bigtable, follow these guidelines as you plan and execute your test:

  • Test with enough data.
    • If the tables in your production instance contain a total of 100 GB of data or less per node, test with a table of the same amount of data.
    • If the tables contain more than 100 GB of data per node, test with a table that contains at least 100 GB of data per node. For example, if your production instance has one four-node cluster, and the tables in the instance contain a total of 1 TB of data, run your test using a table of at least 400 GB.
  • Test with a single table.
  • Stay below the recommended storage utilization per node. For details, see Storage utilization per node.
  • Before you test, run a heavy pre-test for several minutes. This step gives Bigtable a chance to balance data across your nodes based on the access patterns it observes.
  • Run your test for at least 10 minutes. This step lets Bigtable further optimize your data, and it helps ensure that you will test reads from disk as well as cached reads from memory.

If you think that Bigtable might be creating a performance bottleneck in your application, be sure to check all of the following:

  • Look at the Key Visualizer scans for your table. The Key Visualizer tool for Bigtable provides daily scans that show the usage patterns for each table in a cluster. Key Visualizer makes it possible to check whether your usage patterns are causing undesirable results, such as hotspots on specific rows or excessive CPU utilization. Learn how to get started with Key Visualizer.
  • Try commenting out the code that performs Bigtable reads and writes. If the performance issue disappears, then you're probably using Bigtable in a way that results in suboptimal performance. If the performance issue persists, the issue is probably not related to Bigtable.
  • Ensure that you're creating as few clients as possible. Creating a client for Bigtable is a relatively expensive operation. Therefore, you should create the smallest possible number of clients:

    • If you use replication, or if you use app profiles to identify different types of traffic to your instance, create one client per app profile and share the clients throughout your application.
    • If you don't use replication or app profiles, create a single client and share it throughout your application.

    If you're using the HBase client for Java, you create a Connection object rather than a client, so you should create as few connections as possible.

  • Make sure you're reading and writing many different rows in your table. Bigtable performs best when reads and writes are evenly distributed throughout your table, which helps Bigtable distribute the workload across all of the nodes in your cluster. If reads and writes cannot be spread across all of your Bigtable nodes, performance will suffer.

    If you find that you're reading and writing only a small number of rows, you might need to redesign your schema so that reads and writes are more evenly distributed.

  • Verify that you see approximately the same performance for reads and writes. If you find that reads are much faster than writes, you might be trying to read row keys that do not exist, or a large range of row keys that contains only a small number of rows.

    To make a valid comparison between reads and writes, you should aim for at least 90% of your reads to return valid results. Also, if you're reading a large range of row keys, measure performance based on the actual number of rows in that range, rather than the maximum number of rows that could exist.

  • Use the right type of write requests for your data. Choosing the optimal way to write your data helps maintain high performance.

  • Check the latency for a single row. If you observe unexpected latency when sending ReadRows requests, you can check the latency of the first row of the request to narrow down the cause. By default, the overall latency for a ReadRows request includes the latency for every row in the request as well as the processing time between rows. If the overall latency is high but the first row latency is low, this suggests that the latency is caused by the number of requests or processing time, rather than by a problem with Bigtable.

    If you're using the Cloud Bigtable client library for Java, you can view the read_rows_first_row_latency metric in the Cloud Console Metrics Explorer after enabling client side metrics.

BigTable vs HBase: First chapter reading | HBase: The Definitive guide

Oct. 6, 2021

I like to google the following statements:

  1.  block I/O operations ? 
  2. Table scans run in linear time and row key lookups or mutations are performed in logarithmic order—or, in extreme cases, even constant order (using Bloom filters).

Billions of rows * millions of columns * thousands of versions = terabytes or petabytes of storage

We have seen how the Bigtable storage architecture is using many servers to distribute ranges of rows sorted by their key for load-balancing purposes, and can scale to petabytes of data on thousands of machines. The storage format used is ideal for reading adjacent key/value pairs and is optimized for block I/O operations that can saturate disk transfer channels.

Table scans run in linear time and row key lookups or mutations are performed in logarithmic order—or, in extreme cases, even constant order (using Bloom filters).

Designing the schema in a way to completely avoid explicit locking, combined with row-level atomicity, gives you the ability to scale your system without any notable effect on read or write performance.

The column-oriented architecture allows for huge, wide, sparse tables as storing NULLs is free. Because each row is served by exactly one server, HBase is strongly consistent, and using its multi versioning can help you to avoid edit conflicts caused by concurrent decoupled processes or retain a history of changes.

The actual Bigtable has been in production at Google since at least 2005, and it has been in use for a variety of different use cases, from batch-oriented processing to real time data-serving. The stored data varies from very small (like URLs) to quite large (e.g., web pages and satellite imagery) and yet successfully provides a flexible, high performance solution for many well-known Google products, such as Google Earth, Google Reader, Google Finance, and Google Analytics.

HBase BlockCache 101 | Nick Dimiduk | Crash course in 20 minutes | Book reading: HBase definitive guide

HBase BlockCache 101


 

HBase is a distributed database built around the core concepts of an ordered write log and a log-structured merge tree. As with any database, optimized I/O is a critical concern to HBase. When possible, the priority is to not perform any I/O at all. This means that memory utilization and caching structures are of utmost importance. To this end, HBase maintains two cache structures: the “memory store” and the “block cache”. Memory store, implemented as the MemStore, accumulates data edits as they’re received, buffering them in memory (1). The block cache, an implementation of the BlockCache interface, keeps data blocks resident in memory after they’re read.

The MemStore is important for accessing recent edits. Without the MemStore, accessing that data as it was written into the write log would require reading and deserializing entries back out of that file, at least a O(n)operation. Instead, MemStore maintains a skiplist structure, which enjoys a O(log n) access cost and requires no disk I/O. The MemStore contains just a tiny piece of the data stored in HBase, however.

Servicing reads from the BlockCache is the primary mechanism through which HBase is able to serve random reads with millisecond latency. When a data block is read from HDFS, it is cached in the BlockCache. Subsequent reads of neighboring data – data from the same block – do not suffer the I/O penalty of again retrieving that data from disk (2). It is the BlockCache that will be the remaining focus of this post.

Blocks to cache

Before understanding the BlockCache, it helps to understand what exactly an HBase “block” is. In the HBase context, a block is a single unit of I/O. When writing data out to an HFile, the block is the smallest unit of data written. Likewise, a single block is the smallest amount of data HBase can read back out of an HFile. Be careful not to confuse an HBase block with an HDFS block, or with the blocks of the underlying file system – these are all different (3).

HBase blocks come in 4 varieties: DATA, META, INDEX, and BLOOM.

DATA blocks store user data. When the BLOCKSIZE is specified for a column family, it is a hint for this kind of block. Mind you, it’s only a hint. While flushing the MemStore, HBase will do its best to honor this guideline. After each Cell is written, the writer checks if the amount written is >= the target BLOCKSIZE. If so, it’ll close the current block and start the next one (4).

INDEX and BLOOM blocks serve the same goal; both are used to speed up the read path. INDEX blocks provide an index over the Cells contained in the DATA blocks. BLOOM blocks contain a bloom filter over the same data. The index allows the reader to quickly know where a Cell should be stored. The filter tells the reader when a Cell is definitely absent from the data.

Finally, META blocks store information about the HFile itself and other sundry information – metadata, as you might expect. A more comprehensive overview of the HFile formats and the roles of various block types is provided in Apache HBase I/O – HFile.

HBase BlockCache and its implementations

There is a single BlockCache instance in a region server, which means all data from all regions hosted by that server share the same cache pool (5). The BlockCache is instantiated at region server startup and is retained for the entire lifetime of the process. Traditionally, HBase provided only a single BlockCache implementation: the LruBlockCache. The 0.92 release introduced the first alternative in HBASE-4027: the SlabCache. HBase 0.96 introduced another option via HBASE-7404, called the BucketCache.

The key difference between the tried-and-true LruBlockCache and these alternatives is the way they manage memory. Specifically, LruBlockCache is a data structure that resides entirely on the JVM heap, while the other two are able to take advantage of memory from outside of the JVM heap. This is an important distinction because JVM heap memory is managed by the JVM Garbage Collector, while the others are not. In the cases of SlabCache and BucketCache, the idea is to reduce the GC pressure experienced by the region server process by reducing the number of objects retained on the heap.

LruBlockCache

This is the default implementation. Data blocks are cached in JVM heap using this implementation. It is subdivided into three areas: single-access, multi-access, and in-memory. The areas are sized at 25%, 50%, 25% of the total BlockCache size, respectively (6). A block initially read from HDFS is populated in the single-access area. Consecutive accesses promote that block into the multi-access area. The in-memory area is reserved for blocks loaded from column families flagged as IN_MEMORY. Regardless of area, old blocks are evicted to make room for new blocks using a Least-Recently-Used algorithm, hence the “Lru” in “LruBlockCache”.

SlabCache

This implementation allocates areas of memory outside of the JVM heap using DirectByteBuffers. These areas provide the body of this BlockCache. The precise area in which a particular block will be placed is based on the size of the block. By default, two areas are allocated, consuming 80% and 20% of the total configured off-heap cache size, respectively. The former is used to cache blocks that are approximately the target block size (7). The latter holds blocks that are approximately 2x the target block size. A block is placed into the smallest area where it can fit. If the cache encounters a block larger than can fit in either area, that block will not be cached. Like LruBlockCache, block eviction is managed using an LRU algorithm.

BucketCache

This implementation can be configured to operate in one of three different modes: heap, offheap, and file. Regardless of operating mode, the BucketCache manages areas of memory called “buckets” for holding cached blocks. Each bucket is created with a target block size. The heap implementation creates those buckets on the JVM heap; offheap implementation uses DirectByteByffers to manage buckets outside of the JVM heap; filemode expects a path to a file on the filesystem wherein the buckets are created. file mode is intended for use with a low-latency backing store – an in-memory filesystem, or perhaps a file sitting on SSD storage (8). Regardless of mode, BucketCache creates 14 buckets of different sizes. It uses frequency of block access to inform utilization, just like LruBlockCache, and has the same single-access, multi-access, and in-memory breakdown of 25%, 50%, 25%. Also like the default cache, block eviction is managed using an LRU algorithm.

Multi-Level Caching

Both the SlabCache and BucketCache are designed to be used as part of a multi-level caching strategy. Thus, some portion of the total BlockCache size is allotted to an LruBlockCache instance. This instance acts as the first level cache, “L1,” while the other cache instance is treated as the second level cache, “L2.” However, the interaction between LruBlockCache and SlabCache is different from how the LruBlockCache and the BucketCache interact.

The SlabCache strategy, called DoubleBlockCache, is to always cache blocks in both the L1 and L2 caches. The two cache levels operate independently: both are checked when retrieving a block and each evicts blocks without regard for the other. The BucketCache strategy, called CombinedBlockCache, uses the L1 cache exclusively for Bloom and Index blocks. Data blocks are sent directly to the L2 cache. In the event of L1 block eviction, rather than being discarded entirely, that block is demoted to the L2 cache.

Which to choose?

There are two reasons to consider enabling one of the alternative BlockCache implementations. The first is simply the amount of RAM you can dedicate to the region server. Community wisdom recognizes the upper limit of the JVM heap, as far as the region server is concerned, to be somewhere between 14GB and 31GB (9). The precise limit usually depends on a combination of hardware profile, cluster configuration, the shape of data tables, and application access patterns. You’ll know you’ve entered the danger zone when GC pauses and RegionTooBusyExceptions start flooding your logs.

The other time to consider an alternative cache is when response latency really matters. Keeping the heap down around 8-12GB allows the CMS collector to run very smoothly (10), which has measurable impact on the 99th percentile of response times. Given this restriction, the only choices are to explore an alternative garbage collector or take one of these off-heap implementations for a spin.

This second option is exactly what I’ve done. In my next post, I’ll share some unscientific-but-informative experiment results where I compare the response times for different BlockCache implementations.

As always, stay tuned and keep on with the HBase!


1: The MemStore accumulates data edits as they’re received, buffering them in memory. This serves two purposes: it increases the total amount of data written to disk in a single operation, and it retains those recent changes in memory for subsequent access in the form of low-latency reads. The former is important as it keeps HBase write chunks roughly in sync with HDFS block sizes, aligning HBase access patterns with underlying HDFS storage. The latter is self-explanatory, facilitating read requests to recently written data. It’s worth pointing out that this structure is not involved in data durability. Edits are also written to the ordered write log, the HLog, which involves an HDFS append operation at a configurable interval, usually immediate.

2: Re-reading data from the local file system is the best-case scenario. HDFS is a distributed file system, after all, so the worst case requires reading that block over the network. HBase does its best to maintain data locality. These two articles provide an in-depth look at what data locality means for HBase and how its managed.

3: File system, HDFS, and HBase blocks are all different but related. The modern I/O subsystem is many layers of abstraction on top of abstraction. Core to that abstraction is the concept of a single unit of data, referred to as a “block”. Hence, all three of these storage layers define their own block, each of their own size. In general, a larger block size means increased sequential access throughput. A smaller block size facilitates faster random access.

4: Placing the BLOCKSIZE check after data is written has two ramifications. A single Cell is the smallest unit of data written to a DATA block. It also means a Cell cannot span multiple blocks.

5: This is different from the MemStore, for which there is a separate instance for every region hosted by the region server.

6: Until very recently, these memory partitions were statically defined; there was no way to override the 25/50/25 split. A given segment, the multi-access area for instance, could grow larger than it’s 50% allotment as long as the other areas were under-utilized. Increased utilization in the other areas will evict entries from the multi-access area until the 25/50/25 balance is attained. The operator could not change these default sizes. HBASE-10263, shipping in HBase 0.98.0, introduces configuration parameters for these sizes. The flexible behavior is retained.

7: The “approximately” business is to allow some wiggle room in block sizes. HBase block size is a rough target or hint, not a strictly enforced constraint. The exact size of any particular data block will depend on the target block size and the size of the Cell values contained therein. The block size hint is specified as the default block size of 64kb.

8: Using the BucketCache in file mode with a persistent backing store has another benefit: persistence. On startup, it will look for existing data in the cache and verify its validity.

9: As I understand it, there’s two components advising the upper bound on this range. First is a limit on JVM object addressability. The JVM is able to reference an object on the heap with a 32-bit relative address instead of the full 64-bit native address. This optimization is only possible if the total heap size is less than 32GB. See Compressed Oops for more details. The second is the ability of the garbage collector to keep up with the amount of object churn in the system. From what I can tell, the three sources of object churn are MemStore, BlockCache, and network operations. The first is mitigated by the MemSlab feature, enabled by default. The second is influenced by the size of your dataset vs. the size of the cache. The third cannot be helped so long as HBase makes use of a network stack that relies on data copy.

10: Just like with 8, this is assuming “modern hardware”. The interactions here are quite complex and well beyond the scope of a single blog post.

BigTable vs HBase

Oct. 6, 2021

Introduction

I like to google and get more content about the following ideas:

  1. block cache vs key value cache
  2. HBase and BigTable compression algorithms
  3. HBase cannot map storage files into memory, something that is available in Bigtable - Memtable
  4. locality groups - compression 
  5. Column families in Bigtable are used for accounting and access control
  6. Commit log - learn more about this topic 
  7. Study two statements: 
  8. BigTable: Bigtable can memory-map entire storage files and use them to perform lookups without a single disk seek. 
  9. HBase: HBase has an in-memory option per column family and uses its LRU cache‡ to retain blocks for subsequent use.

Overall, HBase implements close to all of the features described in Chapter 1. Where it differs, it may have to because either the Bigtable paper was not very clear to begin with, or it relies on other open source projects to provide various services and those simply work differently.

HBase stores timestamps in milliseconds—as opposed to Bigtable, which uses microseconds. This is not much of an issue and can possibly be attributed to C and Java having different preferred timer resolutions.

While we have not yet addressed the specific details, it should be pointed out that both also use different compression algorithms. HBase uses those supplied in Java, but can also use LZO (with a bit of work; we will look into this later).* Bigtable has a two-phase compression using BMDiff and Zippy.

HBase has coprocessors that are different from what Sawzall, the scripting language used in Bigtable to filter or aggregate data, or the Bigtable Coprocessor framework,† provides. The details on Google’s coprocessor implementation are rather sketchy, so if there are more differences, they are unknown. On the other hand, HBase has support for server-side filters that help reduce the amount of data being  moved from the server to the client.

HBase does primarily work with the Hadoop Distributed File System (HDFS), while Bigtable uses GFS. But HBase can also work on other filesystems thanks to the pluggable FileSystem class provided by Hadoop. There are implementations for Amazon S3 (raw or emulated HDFS), as well as EBS.

HBase cannot map storage files into memory, something that is available in Bigtable. There is ongoing work in HBase to optimize I/O performance, and with the addition of more widespread use of Java’s New I/O (NIO), it may be something that could be enhanced.

Bigtable has a concept called locality groups, which allow the client to group specific column families together and apply shared features, such as compression. This is also useful when the contained columns are accessed together, as all the data is stored in the same storage files. Column families in Bigtable are used for accounting and access control. In HBase, on the other hand, there is only the concept of column families, combining the features that Bigtable has in two distinct concepts.

Apart from the block cache that both systems have, Bigtable also implements a key/value cache, probably for cells that are accessed a lot. 

The handling and implementation of the commit log also differs slightly. Bigtable has two commit logs to handle slow writes and is able to switch between them to compensate for that. This could be implemented in HBase, but it does not seem to be a topic for discussion, and therefore is omitted for the time being.

In contrast, HBase has an option to skip the commit log completely on writes for performance reasons and when the possibility of not being able to replay those logs after a server crash is acceptable.

The METADATA table in Bigtable is also used to store secondary information such as log events related to each tablet. This historical data can be used to analyze tablet transitions, splits, and/or merges. HBase had the notion of a historian in earlier versions that implemented the same concept, but its performance was not good enough and it has been removed.

While splitting regions/tablets is the same for both, merging is handled differently. HBase has a tool that helps you to merge regions manually, while in Bigtable this is handled automatically by the master. Merging in HBase is a delicate operation and currently is left to the operator to decide what is best.

Another very minor difference is that the master in Bigtable is doing the garbage collection of obsolete storage files. One reason for this could be the fact that, in Bigtable, the storage files are tracked in the METADATA table. For HBase, the cleanup is done by the region server that has done the split and no file location is recorded explicitly.

Bigtable can memory-map entire storage files and use them to perform lookups without a single disk seek. HBase has an in-memory option per column family and uses its LRU cache‡ to retain blocks for subsequent use.

There are also some differences in the compaction algorithms. For example, a merging compaction also includes a memtable flush. Mostly, though, they are the same and simply use different names.

Region names, as stored in the meta table in HBase, are a combination of the table name, the start row key, and an ID. In Bigtable, the corresponding tablet names consist of the table identifier and the end row. This has a few implications when it comes to locating data in the storage files (see “Read Path” on page 342). 

Finally, it can be noted that HBase has two separate catalog tables, -ROOT and .META., while in Bigtable the root table, since in both systems it only ever consists of one single region/tablet, is stored as part of the meta table. The first tablet in the METADATA table is the root tablet, and all subsequent ones are the meta tablets. This is just an implementation detail.



Breakfast and learn: Facebook whistleblower Frances Haugen delivers opening statement

Tuesday, October 5, 2021

Facebook whistleblower Frances Haugen delivers opening statement

Oct. 5, 2021

Here is the link. 

Facebook whistleblower and former Facebook product manager Frances Haugen testifies before a Senate subcommittee about Facebook's handling of data of children and other users. For access to live and exclusive video from CNBC subscribe to CNBC PRO: https://cnb.cx/2NGeIvi Facebook whistleblower Frances Haugen told a Senate panel Tuesday that Congress must intervene to solve the “crisis” created by her former employer’s products. The former Facebook product manager for civic misinformation told lawmakers that Facebook consistently puts its own profits over users’ health and safety, which is largely a result of its algorithms’ design that steers users toward high-engagement posts that in some cases can be more harmful.

Though she stopped short of accusing top executives of intentionally creating harmful products, she said that ultimately CEO Mark Zuckerberg had to be responsible for the impact of his business. Haugen also said that Facebook’s algorithm could steer young users from something relatively innocuous such as healthy recipes to content promoting anorexia in a short period of time. She proposed a solution for Facebook to change its algorithms to stop focusing on delivering posts that create more engagement and instead create a chronological feed of posts for Facebook users. That, she said, would help Facebook deliver safer content.

Haugen, who unmasked herself Sunday as the source behind leaked documents at the core of a revealing Wall Street Journal series about Facebook, testified before the Senate Commerce subcommittee on consumer protection. Haugen told CBS’ “60 Minutes” in an interview aired this weekend that the problems she saw at Facebook were worse than anywhere else she’d worked, which includes Google, Yelp and Pinterest. She told the news program that she copied tens of thousands of pages of internal research that she took with her when she left Facebook in May. “I saw that Facebook repeatedly encountered conflicts between its own profits and our safety,” Haugen said in her written testimony. “Facebook consistently resolved those conflicts in favor of its own profits. The result has been a system that amplifies division, extremism and polarization — and undermining societies around the world.”

In her prepared remarks, Haugen said she believes she did the right thing in coming forward but is aware Facebook could use its immense resources to “destroy” her. “I came forward because I recognized a frightening truth: almost no one outside of Facebook knows what happens inside Facebook,” Haugen said in her written remarks. “The company’s leadership keeps vital information from the public, the U.S. government, its shareholders, and governments around the world.” Haugen said a turning point that convinced her of the need to bring information outside Facebook was when the company dissolved the civic integrity team after the 2020 U.S. election. Facebook said it would integrate those responsibilities into other parts of the company. But Haugen said that within six months of the reorganization, 75% of her “pod” of seven people who had mostly come from civic integrity left for other parts of the company or left entirely. “Six months after the reorganization, we had clearly lost faith that those changes were coming,” she said.

In a statement after the hearing concluded, a Facebook spokesperson attempted to cast doubt on Haugen’s credibility.

Frances Haugens

 因此在这种风头期间,脸书的总瘫痪与对于全球网络的灾难性影响,也重新让全球社会开始讨论起「拆割查克柏格数字帝国」的周边讨论。而此一争辩,也正好与24小时前震动美国社会的「脸书吹哨人事件」有关。


脸者吹哨人Frances Haugens在CBS时事节目「60分钟」 中揭发脸书双重标准。(美联社)

豪根揭发脸书「败德报告」

脸书吹哨人事件指的是Facebook的前项目经理法兰西斯.豪根,因为不满老东家在数字道德政策上的严重瑕疵与恶意破坏,而在今年初离职后,将数份关键的「内部调查报告」外流给华尔街日报,并从9月起刊出一连串重量级爆料的〈脸书报告门〉。豪根所揭发的脸书「内部败德报告」,主要拆分成四大指控:

(1)脸书为了扩张流量影响力而秘密创建「社群阶级」,高知名用户、大型政治用户、高声量用户可以不受社群规则的交叉核实限制。这种鉴于「声量阶级」而非「事实性评断」的差别待遇,是造成全球性社群内战与「谣言系意见领袖」崛起的关键原因。

(2)脸书的持股人与创办人查克柏格存在严重分裂,甚至因剑桥分析丑闻的美国政府和解禁问题对簿公堂,股东认为查克柏格把脸书当成私产,并用公司资源来保护自己不受法律制裁。

(3)脸书的内部研究发现「Instagram会对青少女族群带来特别严重的身心压力」,相关因素包括网络性骚扰、匿名霸凌与外貌羞辱等,负面身心回馈比例高达32%——但对此,脸书的作法是「压制这份研究的曝光」,并完全不采取任何改善行动。

(4)在「公义之前Facebook选择『无限赚钱』」,虽然获利是公司的天职,但Facebook偏执的盈利导向却在言论审查、网络暴力、用户隐私、儿少身心健康、社会冲突与假新闻之间造成全球性严重伤害——虽然脸书总是声称会做事实查核与不当言论管理,「可是只要选举结束,风头一过,Facebook就会把这些『防火墙机制』秘密关闭——将公司获利成长至于整体安全之前,这就像是对民主社会的背叛。

脸书瘫痪的破坏程度,已严重阻碍了全球网络的数字服务运作。(欧新社)

豪根的爆料在华尔街日报刊出后,引发了国际舆论的热烈讨论,而他本人也在10月3日的CBS专访节目中以「真人现身」的方式亮相。尽管对于豪根的控诉,脸书总部已发出严重的谴责与驳斥不实声明,但豪根原本就将于10月5日出席的参议院社群问题听证会,此时再配上脸书大崩溃的全球冲击,「大到不能倒」的社群时代争辩,料也将因此风波而更为升级。

Tuesday, October 5, 2021 Oct. 4, 2021 Facebook outage | Chinese article

 全球网民依赖的社群平台脸书与相关平台Instagram及WhatsApp,4日上午突然当机,近6小时后才恢复正常。(路透)


「Faceook的我们正试着全力排除这场『不明故障』,抱歉...via Twitter。」全球社群龙头Facebook(脸书)集团自美东时间4日早上11时39分起,传出极为严重的「全集团瘫痪」——包括Facebook、脸书 Messenger、Instagram、Whatsapp与脸书投资的VR系统Oculus,瞬间于全球互联网大断线。

相关障害一直到晚上6时左右才排除。但Facebook全家族无预警下线的6个半小时内,却严重阻断了全球30亿用户的日常通信,甚至直接阻断了其他网络业者的帐号服务。故障期间,Facebook的全球市值一口气蒸发200到250亿美元,恐是脸书系列集体下线最久、瘫痪损害规模最大的「严重故障事件」。

集团内部 抢修应变都失灵

根据华尔街日报所取得的脸书集团内部通知信,截至4日中午为止,Facebook总公司都无法向员工说明瘫痪原因与公司到底发生了什么事,集团内部的抢修、紧急应变与客户通知系统也完全失灵,仅强调就已知资讯而言「本回故障并不似『恶意事件』」。其他外部专家与观察组织则认为,Facebook所遭遇的可能是域名系统(DNS)的内部设置故障,造成全球互联网与脸书集团服务的IP「彻底断线」。

Facebook的大断线,同时也让Twitter、Telegram、甚至是中国的Tiktok,见猎心喜地大开嘲讽并趁势抢进服务。但这起耗时超过6小时的「无端总崩溃」,却也可能加剧欧美政坛对于脸书创办人查克柏格(Mark Zuckerberg)的「社群帝国」疑虑——特别是在脸书总崩的24小时前,近一个月在华尔街日报揭发〈脸书报告〉的前脸书员工吹哨人豪根(Frances Haugen)才公开现身,全力控诉脸书集团在查克柏格治下已经「明知有恶而硬要为之」,对全球政治、社会包容、用户身心健康极具风险与威胁的方式自顾获利。

安全、道德、垄断都遭质疑

由于豪根本人即将于美东时间5日,参加美国国会参议院针对Facebook管制问题的重要听证会,关键交锋前夕的「脸书全球大崩溃」,也重新引发了世界各地对于脸书「大到不能倒」的安全、道德与过度垄断之质疑?

全球网民依赖的社群平台脸书与相关平台Instagram及WhatsApp,4日上午突然当机, 脸书用户重批脸书,在Twitter上极尽挖苦之能事。(取材自脸书)

金融时报表示:在查克柏格帝国「全球下线」的6小时内,脸书总部一直无法有效厘清、或对公众更新「障害排除进度」。但瘫痪的破坏程度,却已严重阻碍了全球网络的数字服务运作,因为除了Facebook、Instagram的「一般社群功能」失效外,数十亿个「连动Facebook帐号登录」的其他网络帐号(包括其他App、游戏、媒体订阅、电商服务...)也都一并无法使用。

社群平台脸书与相关平台Instagram及WhatsApp,4日上午突然当机。(欧新社)

再加上流行于世界各地的即时通信软件WhatsApp,也已由脸书集团持有,一起爆炸的状况也让超过20亿通信用户「与日常失联」——因此,脸书集团的总崩溃不仅是单纯「社群下线」,以全球规模损害并影响的社群广告、电子商务、数字帐号、甚至日常通信,都造成了极大而罕有的震荡冲击。

事实上,在过去Facebook集团就是常出现「集体性故障」的问题,不过受影响的区域大多有限,比较少出现「整个单一集团总崩盘」,并在完全不清楚也不通知的状况下,以「全球范围」总故障竟持续6小时以上的长时间。

「这真的很扯。」

Kentik网络分析主任马多利(Doug Madory)就对华盛顿邮报的访问如此表示:「Facebook的所有服务完全爆掉了。」

华尔街日报报导,Facebook的总崩溃受害范围遍及「脸书集团的所有持有服务」。在一封被《华尔街日报》取得的员工内部信里,脸书总公司更是直接表示:「我们不知道发生了什么事,目前工程抢修团队正紧急前往加州机房总部试图『物理还原』。

全球网民依赖的社群平台脸书与相关平台Instagram及WhatsApp,4日上午突然当机,近6小时后才恢复正常。(路透)

抢修进度 尴尬用推特发布

相关障害,一直持续到晚上6时才突然「逐步恢复」。但在这一过程中,脸书集团的所有公开道歉信与抢修进度说明,却得尴尬地使用「社群竞媒」Twitter来对外发布,脸书Messegenger的Twitter帐号甚至尴尬地试着以冷笑话回应呛声的用户:

「水星逆行的威力真的很强呢!」

但对此,大多数的商业用户与网络观察者却都笑不出来。因为脸书的无预警总崩,不仅在短短6小时内,对自家集团带来200亿美元以上的损失,通信的大乱还从脸书本体扩张到「所有链接Facebook帐号的网络服务」;与此同时,虽然Twitter、Telegram等社群软件不受影响,但从Facebook瘫痪而「瞬间外流」的数十亿用户流量,却也短暂造成了各路跨国网站的不稳与震荡,甚至从中加剧了部分谣言所煽动的「网战总开战恐慌」。

脸书集团史上最大的总瘫痪,在国际舆论上也重新燃起「查克柏格帝国是否大到不能倒?」的政策激辩——此一讨论不仅只限于欧洲的数字反垄断政策 ,在美国政坛也因为川普被停权禁言、以及1月6日国会山庄遇袭事件,而正在各种针对Facebook的听证会车轮战。

全球网民依赖的社群平台脸书与相关平台Instagram及WhatsApp,4日上午突然当机,近6小时后才恢复正常。(路透)





Oct. 4, 2021 Facebook outage | Chinese article

 当地时间 4 日,美国互联网巨头脸书公司的多款社交软件,在全球多个国家出现中断,甚至瘫痪,持续了近 7 个小时。据法新社报道,这是脸书公司成立以来最严重的一次宕机故障。

据报道,此次宕机故障涉及脸书旗下的四大社交产品,包括脸书平台、即时通讯软件移动聊天服务 Messenger、WhatsApp,以及图片社交服务 Instagram。随后,随后谷歌、亚马逊等网站也出现大面积宕机情况。近 7 个小时后,脸书表示系统已恢复运行。

截至目前,此次服务器大规模宕机的确切原因尚不清楚。脸书网页上显示服务中断是由于域名系统错误。有业内人士表示,没有证据表明这次系统中断是一次网络攻击,最有可能的解释是脸书在维护过程中错误地关闭了网络。

安克诺斯数据保护公司副总裁 斯特:由于某种原因,脸书的域名系统服务器似乎从互联网上消失了,无法经路由中转,通过域名进行访问,糟糕的是,故障也影响到脸书的技术人员,现在可能必须去现场进行修复。

《纽约时报》一名叫弗兰克尔的记者 4 日在推特上表示,脸书员工无法进入公司大楼评估社交网络运行故障,因为他们的工牌无法打开门锁。脸书创始人马克 · 扎克伯格当天也因宕机故障公开致歉。

受此影响,脸书股价 4 日暴跌近 5%,创近一年以来最大单日跌幅,6 百亿多美元,约合人民币 4000 多亿元的市值瞬间蒸发。