0% encontró este documento útil (0 votos)
6 vistas26 páginas

Bigtable: Almacenamiento Distribuido Eficaz

Bigtable es un sistema de almacenamiento distribuido diseñado por Google para gestionar datos estructurados a gran escala, capaz de manejar petabytes de información en miles de servidores. Utilizado por diversos proyectos de Google, ofrece un modelo de datos flexible y un alto rendimiento, permitiendo a los usuarios controlar el diseño y formato de sus datos. El artículo detalla el modelo de datos, la API y la implementación de Bigtable, así como ejemplos de su uso en diferentes aplicaciones.

Cargado por

OS KRG
Derechos de autor
© All Rights Reserved
Nos tomamos en serio los derechos de los contenidos. Si sospechas que se trata de tu contenido, reclámalo aquí.
Formatos disponibles
Descarga como PDF, TXT o lee en línea desde Scribd
0% encontró este documento útil (0 votos)
6 vistas26 páginas

Bigtable: Almacenamiento Distribuido Eficaz

Bigtable es un sistema de almacenamiento distribuido diseñado por Google para gestionar datos estructurados a gran escala, capaz de manejar petabytes de información en miles de servidores. Utilizado por diversos proyectos de Google, ofrece un modelo de datos flexible y un alto rendimiento, permitiendo a los usuarios controlar el diseño y formato de sus datos. El artículo detalla el modelo de datos, la API y la implementación de Bigtable, así como ejemplos de su uso en diferentes aplicaciones.

Cargado por

OS KRG
Derechos de autor
© All Rights Reserved
Nos tomamos en serio los derechos de los contenidos. Si sospechas que se trata de tu contenido, reclámalo aquí.
Formatos disponibles
Descarga como PDF, TXT o lee en línea desde Scribd

Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido


para datos estructurados

FAY CHANG, JEFFREY DEAN, SANJAY GHEMAWAT, WILSON C. HSIEH,


DEBORAH A. WALLACH, MIKE BURROWS, TUSHAR CHANDRA,
ANDREW FIKES y ROBERT E. GRUBER
Google, Inc.

Bigtable es un sistema de almacenamiento distribuido para la gestión de datos estructurados, diseñado para escalar a un tamaño muy
grande: petabytes de datos en miles de servidores de servicios básicos. Muchos proyectos de Google almacenan datos en Bigtable,
como la indexación web, Google Earth y Google Finance. Estas aplicaciones imponen exigencias muy diferentes a Bigtable, tanto en
términos de tamaño de datos (desde URL hasta páginas web e imágenes satelitales) como de latencia (desde el procesamiento masivo
de backend hasta la entrega de datos en tiempo real). A pesar de estas diversas exigencias, Bigtable ha proporcionado con éxito una
solución flexible y de alto rendimiento para todos estos productos de Google. En este artículo, describimos el sencillo modelo de datos
que ofrece Bigtable, que ofrece a los clientes un control dinámico sobre el diseño y el formato de los datos, y describimos el diseño y la
implementación de Bigtable.

Categorías y descriptores de temas: C.2.4 [Redes de comunicación informática]: Sistemas distribuidos: bases de datos distribuidas

Condiciones generales: Diseño

Palabras y frases clave adicionales: Almacenamiento distribuido a gran escala

Formato de Referencia ACM: Chang,


F., Dean, J., Ghemawat, S., Hsieh, WC, Wallach, DA, Burrows, M., Chandra, T., Fikes, A. y Gruber, RE 2008. Bigtable: Un sistema de
almacenamiento distribuido para datos estructurados. ACM Trans. Comput. Syst. 26, 2, Artículo 4 (junio de 2008), 26 páginas. DOI =
10.1145/1365815.1365816. [Link]

Este artículo se publicó originalmente como trabajo premiado en las Actas del 7.º Simposio sobre Diseño e Implementación de Sistemas
Operativos [Chang et al., 2006]. Se republica aquí con pequeñas modificaciones y aclaraciones.

Dirección de los autores: Google Inc., 1600 Amphitheatre Parkway, Mountain View, CA 94043; correo electrónico: {fay, jeff, sanjay,
wilsonh, kerr, m3b, tushar, fikes, gruber}@[Link].
Se concede permiso para realizar copias digitales o impresas de parte o la totalidad de esta obra para uso personal o académico sin
costo alguno, siempre que no se realicen ni distribuyan con fines de lucro o para obtener una ventaja comercial directa, y que las copias
muestren este aviso en la primera página o pantalla inicial de una pantalla, junto con la cita completa. Se deben respetar los derechos
de autor de los componentes de esta obra que sean propiedad de terceros. Se permite la inclusión de resúmenes con créditos. Para
copiar, republicar, publicar en servidores, redistribuir a listas o utilizar cualquier componente de esta obra en otras obras, se requiere un
permiso previo específico o el pago de una tarifa. Los permisos pueden solicitarse al Departamento de Publicaciones, ACM, Inc., 2
Penn Plaza, Suite 701, Nueva York, NY 10121­0701, EE. UU., fax +1 (212) 869­0481 o permission@[Link]. c 2008 ACM
0734­2071/2008/06­ART4 $5,00 DOI: 10.1145/1365815.1365816. [Link]

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:2 ∙ F. Chang y otros.

1. INTRODUCCIÓN

A finales de 2003, diseñamos, implementamos y desplegamos un sistema de almacenamiento distribuido


para administrar datos estructurados en Google llamado Bigtable.
Bigtable está diseñado para escalar de forma fiable a petabytes de datos y miles de máquinas. Bigtable
ha logrado varios objetivos: amplia aplicabilidad, escalabilidad y alta...
Rendimiento y alta disponibilidad. Bigtable es utilizado por más de sesenta empresas de Google.
Productos y proyectos, incluidos Google Analytics, Google Finance, Orkut, Búsqueda personalizada,
Writely y Google Earth. Estos productos utilizan Bigtable para...
Variedad de cargas de trabajo exigentes, que van desde trabajos de procesamiento por lotes orientados
al rendimiento hasta la entrega de datos a usuarios finales con alta latencia. Bigtable
Los clústeres utilizados por estos productos abarcan una amplia gama de configuraciones, desde una
desde unos pocos hasta miles de servidores y almacenar hasta varios cientos de terabytes
de datos.
En muchos sentidos, Bigtable se asemeja a una base de datos: comparte muchas estrategias de
implementación con las bases de datos. Bases de datos paralelas [DeWitt y Gray, 1992]
y las bases de datos de memoria principal [DeWitt et al. 1984] han logrado escalabilidad
y de alto rendimiento, pero Bigtable proporciona una interfaz diferente a la de estos.
Bigtable no admite un modelo de datos relacional completo; en cambio, proporciona a los clientes un
modelo de datos simple que permite el control dinámico de los datos.
diseño y formato, y permite a los clientes razonar sobre las propiedades de localidad de
Los datos representados en el almacenamiento subyacente. Los datos se indexan mediante filas y
Nombres de columna que pueden ser cadenas arbitrarias. Bigtable también trata los datos como cadenas
sin interpretar, aunque los clientes suelen serializar diversas formas de datos estructurados.
y datos semiestructurados en estas cadenas. Los clientes pueden controlar la ubicación de
sus datos mediante elecciones cuidadosas en sus esquemas. Finalmente, el esquema de Bigtable...
Los parámetros permiten a los clientes controlar dinámicamente si servir datos desde la memoria o desde
el disco.
La Sección 2 describe el modelo de datos con más detalle y la Sección 3 proporciona una
Descripción general de la API del cliente. La sección 4 describe brevemente la API subyacente de Google.
Infraestructura de la que depende Bigtable. La Sección 5 describe los fundamentos de la implementación
de Bigtable, y la Sección 6 describe algunas de las mejoras que realizamos para mejorar el rendimiento
de Bigtable. La Sección 7 proporciona
Mediciones del rendimiento de Bigtable. Describimos varios ejemplos de cómo
Bigtable se utiliza en Google en la Sección 8 y analizamos algunas lecciones que aprendimos en
Diseño y soporte de Bigtable en la Sección 9. Finalmente, la Sección 10 describe
trabajo relacionado, y la Sección 11 presenta nuestras conclusiones.

2. MODELO DE DATOS

Un clúster de Bigtable es un conjunto de procesos que ejecutan el software de Bigtable. Cada uno
Un clúster sirve para un conjunto de tablas. Una tabla en Bigtable es un mapa ordenado multidimensional,
disperso, distribuido y persistente. Los datos se organizan en tres dimensiones:
filas, columnas y marcas de tiempo.

(fila:cadena, columna:cadena, tiempo:int64) → cadena

Nos referimos al almacenamiento referenciado por una clave de fila, una clave de columna y
Marca de tiempo como una celda. Las filas se agrupan para formar la unidad de carga.
ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido para datos estructurados ∙ 4:3

Fig. 1. Un fragmento de una tabla de ejemplo que almacena páginas web. El nombre de la fila es una URL invertida.
La familia de columnas de contenido contiene el contenido de la página, y la familia de columnas de ancla contiene
el texto de cualquier ancla que haga referencia a la página. La página principal de CNN es referenciada tanto por las
páginas principales de Sports Illustrated como de MY­look, por lo que la fila contiene las columnas llamadas
ancla:[Link] y ancla:[Link]. Cada celda de ancla tiene una versión; la columna de contenido tiene tres
versiones, en las marcas de tiempo t3, t5 y t6.

equilibrio y las columnas se agrupan para formar la unidad de control de acceso y contabilidad de recursos.

Nos decidimos por este modelo de datos tras examinar diversos usos potenciales de un sistema
similar a Bigtable. Consideremos un ejemplo concreto que inspiró muchas de nuestras decisiones de
diseño: una copia de una gran colección de páginas web e información relacionada, que podría utilizarse
en diversos proyectos. Llamaremos a esta tabla Webtable. En Webtable, usaríamos las URL como claves
de fila, diversos aspectos de las páginas web como nombres de columna y almacenaríamos el contenido
de las páginas web en la columna "contents: ", bajo las marcas de tiempo de su obtención, como se
ilustra en la Figura 1.

Filas. Bigtable mantiene los datos en orden lexicográfico por clave de fila. Las claves de fila de una
tabla son cadenas arbitrarias (actualmente de hasta 64 KB, aunque el tamaño típico para la mayoría de
nuestros usuarios es de 10 a 100 bytes). Cada lectura o escritura de datos bajo una sola clave de fila es
serializable (independientemente del número de columnas diferentes que se lean o escriban en la fila),
una decisión de diseño que facilita a los clientes comprender el comportamiento del sistema ante
actualizaciones simultáneas de la misma fila. En otras palabras, la fila es la unidad de consistencia
transaccional en Bigtable, que actualmente no admite transacciones entre filas.

Las filas con claves consecutivas se agrupan en tabletas, que constituyen la unidad de distribución y
balanceo de carga. Como resultado, las lecturas de rangos cortos de filas son eficientes y, por lo general,
requieren comunicación con un número reducido de máquinas. Los clientes pueden aprovechar esta
propiedad seleccionando sus claves de fila para obtener una buena localidad en sus accesos a los datos.
Por ejemplo, en Webtable, las páginas del mismo dominio se agrupan en filas contiguas invirtiendo los
componentes del nombre de host de las URL. Almacenaríamos los datos de [Link]/[Link]
bajo la clave [Link]/[Link]. Almacenar páginas del mismo dominio cerca unas de otras
aumenta la eficiencia de algunos análisis de host y dominio.

Columnas. Las claves de columna se agrupan en conjuntos denominados familias de columnas, que
constituyen la unidad de control de acceso. Todos los datos almacenados en una familia de columnas
suelen ser del mismo tipo (comprimimos los datos de la misma familia de columnas). Una familia de
columnas debe crearse explícitamente antes de poder almacenar datos en cualquier...
ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:4
∙ F. Chang y otros.

Clave de columna de esa familia. Tras crear una familia, se puede usar cualquier clave de columna
dentro de ella: los datos se pueden almacenar bajo dicha clave sin afectar el esquema de la tabla.
Nuestro objetivo es que el número de familias de columnas distintas en una tabla sea pequeño (cientos
como máximo) y que las familias cambien con poca frecuencia durante la operación; esta limitación
evita que los metadatos ampliamente compartidos sean demasiado grandes. Por el contrario, una
tabla puede tener un número ilimitado de columnas.

Se pueden eliminar familias de columnas enteras modificando el esquema de una tabla, en cuyo
caso se eliminan los datos almacenados bajo cualquier clave de columna de esa familia.
Sin embargo, dado que Bigtable no admite transacciones en varias filas, los datos almacenados bajo
una clave de columna particular no se pueden eliminar de forma atómica si residen en varias filas.

Una clave de columna se nombra con la siguiente sintaxis: familia:calificador. Los nombres de las
familias de columnas deben ser imprimibles, pero los calificadores pueden ser cadenas arbitrarias. Un
ejemplo de familia de columnas para Webtable es idioma, que almacena el idioma en el que se
escribió una página web. Usamos solo una clave de columna con un calificador vacío en la familia
idioma para almacenar el ID de idioma de cada página web. Otra familia de columnas útil para esta
tabla es ancla; cada clave de columna de esta familia representa un ancla, como se muestra en la
Figura 1. El calificador es el nombre del sitio web de referencia; la celda contiene el texto asociado al
enlace.
El control de acceso y la contabilidad de disco y memoria se realizan a nivel de familia de columnas.
En nuestro ejemplo de tabla web, estos controles nos permiten gestionar varios tipos de aplicaciones:
algunas que añaden nuevos datos base, otras que los leen y crean familias de columnas derivadas, y
otras que solo pueden ver los datos existentes (y posiblemente ni siquiera todas las familias existentes
por motivos de privacidad).

Marcas de tiempo. Las distintas celdas de una tabla pueden contener varias versiones de los
mismos datos, indexadas por marca de tiempo. Las marcas de tiempo de Bigtable son enteros de 64
bits. Bigtable puede asignarlas implícitamente, representando así el tiempo real en microsegundos, o
las aplicaciones cliente pueden asignarlas explícitamente. Las aplicaciones que necesitan evitar
colisiones deben generar sus propias marcas de tiempo únicas. Las distintas versiones de una celda
se almacenan en orden decreciente de marca de tiempo, de modo que las versiones más recientes se
lean primero.
Para simplificar la gestión de datos versionados, ofrecemos dos configuraciones por familia de
columnas que indican a Bigtable que recolecte automáticamente los datos versionados. El cliente
puede especificar que solo se conserven las últimas n versiones de los datos o que solo se conserven
las versiones suficientemente recientes (por ejemplo, solo los valores escritos en los últimos siete días).

En nuestro ejemplo de Webtable, podemos establecer las marcas de tiempo de las páginas
rastreadas almacenadas en la columna "contents: " con las horas en que se rastrearon realmente
estas versiones. El mecanismo de recolección de elementos no utilizados descrito anteriormente nos
permite indicar a Bigtable que conserve solo las tres versiones más recientes de cada página.

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido para datos estructurados ∙ 4:5

Fig. 2. Escribiendo en Bigtable.

Fig. 3. Lectura desde Bigtable.

3. API

La API de Bigtable proporciona funciones para crear y eliminar tablas y familias de columnas. También
permite modificar metadatos de clústeres, tablas y familias de columnas, como los derechos de control
de acceso.
Las aplicaciones cliente pueden escribir o eliminar valores en Bigtable, buscar valores en filas
individuales o iterar sobre un subconjunto de datos de una tabla. La Figura 2 muestra código C++ que
utiliza una abstracción RowMutation para realizar una serie de actualizaciones. (Se omitieron detalles
irrelevantes para abreviar el ejemplo). La llamada a Apply realiza una mutación atómica en la Webtable:
añade un ancla a [Link] y elimina otro ancla.

La Figura 3 muestra código C++ que utiliza una abstracción de Scanner para iterar sobre todos los
anclajes en una fila específica. Los clientes pueden iterar sobre múltiples familias de columnas, y
existen varios mecanismos para limitar las filas, columnas y marcas de tiempo que recorre un escaneo.
Por ejemplo, podríamos restringir el escaneo anterior para que solo produzca anclajes cuyas columnas
coincidan con la expresión regular anchor:*.[Link], o para que solo produzca anclajes cuyas marcas
de tiempo estén dentro de los diez días posteriores a la hora actual.

Bigtable admite otras funciones que permiten al usuario manipular datos de forma más compleja.
En primer lugar, Bigtable admite transacciones de una sola fila.

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:6 ∙ F. Chang y otros.

Fig. 4. Un conjunto típico de procesos que se ejecutan en una máquina de Google. Una máquina suele ejecutar
numerosos trabajos de distintos usuarios.

que puede utilizarse para realizar secuencias atómicas de lectura, modificación y escritura en
datos almacenados bajo una única clave de fila. Bigtable no admite actualmente transacciones
generales entre claves de fila, aunque proporciona una interfaz para la escritura por lotes en
claves de fila en los clientes. En segundo lugar, Bigtable permite utilizar celdas como
contadores de enteros. Finalmente, Bigtable admite la ejecución de scripts proporcionados
por el cliente en los espacios de direcciones de los servidores. Los scripts se escriben en un
lenguaje llamado Sawzall [Pike et al., 2005], desarrollado en Google para el procesamiento
de datos. Actualmente, nuestra API basada en Sawzall no permite que los scripts del cliente
vuelvan a escribir en Bigtable, pero sí permite diversas formas de transformación de datos,
filtrado basado en expresiones arbitrarias y resumen mediante diversos operadores.

Bigtable se puede usar con MapReduce [Dean y Ghemawat 2004], un framework para
ejecutar cálculos paralelos a gran escala desarrollado en Google. Hemos desarrollado un
conjunto de envoltorios que permiten usar Bigtable como fuente de entrada y como destino
de salida para trabajos de MapReduce.

4. BLOQUES DE CONSTRUCCIÓN

Bigtable se basa en varios otros componentes de la infraestructura de Google. Un clúster de


Bigtable suele operar en un conjunto compartido de máquinas que ejecutan diversas
aplicaciones distribuidas. Bigtable depende de un sistema de gestión de clústeres de Google
para programar tareas, administrar recursos en máquinas compartidas, supervisar el estado
de las máquinas y gestionar fallos. Los procesos de Bigtable suelen compartir las mismas
máquinas con procesos de otras aplicaciones. Por ejemplo, como se ilustra en la Figura 4,
un servidor de Bigtable puede ejecutarse en la misma máquina que un trabajador de
MapReduce [Dean y Ghemawat, 2004], servidores de aplicaciones y un servidor para el
Sistema de Archivos de Google (GFS) [Ghemawat et al., 2003].
ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido para datos estructurados ∙ 4:7

Bigtable utiliza GFS para almacenar archivos de registro y datos. GFS es un sistema de archivos distribuido
que mantiene múltiples réplicas de cada archivo para una mayor confiabilidad y
disponibilidad.
El formato de archivo inmutable SSTable de Google se utiliza internamente para almacenar
Archivos de datos de Bigtable. Una SSTable proporciona un mapa persistente, ordenado e inmutable.
de claves a valores, donde tanto las claves como los valores son cadenas de bytes arbitrarias.
Se proporcionan operaciones para buscar el valor asociado con una clave específica.
y para iterar sobre todos los pares clave/valor en un rango de claves especificado. Internamente, cada
SSTable contiene una secuencia de bloques (de manera predeterminada, cada bloque tiene un tamaño de 64 KB,
Pero el tamaño es configurable. Un índice de bloque (almacenado al final de la SSTable).
se utiliza para localizar bloques; el índice se carga en la memoria cuando se crea la SSTable.
abierto. Se puede realizar una búsqueda con una sola búsqueda en el disco: primero encontramos el
bloque apropiado realizando una búsqueda binaria en el índice en memoria, y
Luego, lea el bloque correspondiente del disco. Opcionalmente, una SSTable puede mapearse completamente
en memoria, lo que permite realizar búsquedas y escaneos.
sin tocar el disco.
Bigtable se basa en un servicio de bloqueo distribuido persistente y de alta disponibilidad
llamado Chubby [Burrows 2006]. Un servicio Chubby consta de cinco réplicas activas, una de las cuales se elige
como maestra y atiende activamente las solicitudes.
El servicio está activo cuando la mayoría de las réplicas están en ejecución y pueden comunicarse entre sí.
Chubby utiliza el algoritmo Paxos [Chandra et al., 2007;
Lamport 1998] para mantener la consistencia de sus réplicas ante fallos. Chubby proporciona un espacio de
nombres compuesto por directorios y archivos pequeños. Cada directorio o
El archivo se puede usar como un bloqueo, y las lecturas y escrituras en un archivo son atómicas. El Chubby
La biblioteca de cliente proporciona un almacenamiento en caché consistente de los archivos de Chubby. Cada cliente de Chubby...
Mantiene una sesión con un servicio de Chubby. La sesión de un cliente expira si no puede renovar su contrato
de sesión dentro del plazo de vencimiento. Cuando un cliente...
La sesión expira, pierde los bloqueos y los identificadores abiertos. Los clientes de Chubby también pueden
Registra devoluciones de llamadas en archivos y directorios de Chubby para notificar cambios o
expiración de la sesión.
Bigtable usa Chubby para una variedad de tareas: para garantizar que haya como máximo
un maestro activo en cualquier momento; para almacenar la ubicación de arranque de los datos de Bigtable
(ver Sección 5.1); para descubrir servidores de tabletas y finalizar las muertes de servidores de tabletas
(véase la Sección 5.2); y para almacenar esquemas de Bigtable (véase la Sección 5.5). Si Chubby
Si deja de estar disponible durante un período prolongado, Bigtable deja de estar disponible. En agosto de 2006,
medimos este efecto en 14 clústeres de Bigtable que abarcaban
11 instancias de Chubby. El porcentaje promedio de horas de servidor de Bigtable durante
El porcentaje de datos almacenados en Bigtable que no estaban disponibles debido a la falta de disponibilidad
de Chubby (causada por interrupciones de Chubby o problemas de red) fue del 0,0047 %.
El porcentaje del grupo individual que fue más afectado por la falta de disponibilidad de Chubby fue del 0,0326%.

5. IMPLEMENTACIÓN

La implementación de Bigtable tiene tres componentes principales: una biblioteca que es


Conectado a cada cliente, un servidor maestro y muchos servidores de tableta. Tableta
Los servidores se pueden agregar (o eliminar) dinámicamente de un clúster para adaptarse
cambios en las cargas de trabajo.

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:8 ∙ F. Chang y otros.

Fig. 5. Jerarquía de ubicación de la tableta.

El maestro se encarga de asignar tabletas a servidores de tabletas, detectar la adición y el vencimiento


de servidores de tabletas, equilibrar la carga entre ambos servidores y recolectar archivos no utilizados en
GFS. Además, gestiona cambios de esquema, como la creación y eliminación de familias de tablas y
columnas.
Cada servidor de tabletas administra un conjunto de tabletas (normalmente tenemos entre diez y mil
tabletas por servidor). El servidor gestiona las solicitudes de lectura y escritura a las tabletas que ha
cargado y también divide las tabletas que han crecido demasiado.

Al igual que con muchos sistemas de almacenamiento distribuido de un solo maestro [Ghemawat et al.
[2003; Hartman y Ousterhout 1993], los datos del cliente no se mueven a través del maestro: los clientes
se comunican directamente con los servidores de tabletas para lecturas y escrituras.
Dado que los clientes de Bigtable tampoco dependen del maestro para obtener información sobre la
ubicación de la tableta, la mayoría nunca se comunica con él. Por lo tanto, en la práctica, el maestro tiene
poca carga.
Un clúster de Bigtable almacena varias tablas. Cada tabla consta de un conjunto de tabletas, y cada
tableta contiene todos los datos asociados a un rango de filas.
Inicialmente, cada tabla consta de una sola tableta. A medida que crece, se divide automáticamente en
varias tabletas, cada una con un tamaño predeterminado de aproximadamente 1 GB.
Aunque nuestro modelo admite datos de cualquier tamaño, la implementación actual de Bigtable no
admite valores extremadamente grandes. Dado que una tableta no se puede dividir en mitad de una fila,
recomendamos a los usuarios que cada fila no contenga más de unos pocos cientos de GB de datos.

En el resto de esta sección describimos algunos de los detalles de la


Implementación de Bigtable.

5.1 Ubicación de la tableta

Utilizamos una jerarquía de tres niveles análoga a la de un árbol B+ [Comer, 1979] para almacenar la
información de ubicación de las tabletas (Figura 5). El primer nivel es un archivo almacenado en Chubby
que contiene la ubicación de la tableta raíz. Esta tableta raíz contiene las ubicaciones de todas las tabletas
de una tabla METADATA especial . Cada tableta METADATA contiene la ubicación de un conjunto de
tabletas de usuario. La tableta raíz se trata
ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido para datos estructurados ∙ 4:9

especialmente (nunca se divide) para garantizar que la jerarquía de ubicación de la tableta no tenga
más de tres niveles.
La tabla METADATA almacena la ubicación de una tableta mediante una clave de fila que codifica
el identificador de tabla de la tableta y su fila final. Cada fila de METADATA almacena aproximadamente
1 KB de datos en memoria. Con un límite moderado de tabletas de METADATA de 128 MB , nuestro
esquema de ubicación de tres niveles es suficiente para direccionar 234 tabletas (o 261 bytes en
tabletas de 128 MB).
La biblioteca cliente recorre la jerarquía de ubicaciones para localizar tabletas y almacena en caché
las ubicaciones que encuentra. Si el cliente desconoce la ubicación de una tableta o descubre que la
información de ubicación almacenada en caché es incorrecta, asciende recursivamente en la
jerarquía. Si la caché del cliente está vacía, el algoritmo de ubicación requiere tres viajes de ida y
vuelta a la red, incluyendo una lectura desde Chubby. Si la caché del cliente está obsoleta, el algoritmo
podría tardar hasta seis viajes de ida y vuelta, ya que las entradas de caché obsoletas solo se detectan
en caso de fallos (lo cual esperamos que sea poco frecuente, ya que las tabletas METADATA no
deberían moverse con mucha frecuencia). Aunque las ubicaciones de las tabletas se almacenan en
memoria, por lo que no se requieren accesos GFS, reducimos aún más este coste en el caso común
al permitir que la biblioteca cliente precapture las ubicaciones de las tabletas: lee los metadatos de
más de una tableta cada vez que lee la tabla METADATA .

También almacenamos información secundaria en la tabla METADATA , incluyendo un registro de


todos los eventos relacionados con cada tableta (por ejemplo, cuándo un servidor comienza a servirla).
Esta información es útil para la depuración y el análisis de rendimiento.

5.2 Asignación de Tabletas.

Cada tableta se asigna a un máximo de un servidor de tabletas a la vez. El maestro realiza un


seguimiento del conjunto de servidores de tabletas activos y de la asignación actual de tabletas a
dichos servidores, incluyendo las tabletas sin asignar. Cuando una tableta no está asignada y hay un
servidor de tabletas con espacio suficiente, el maestro la asigna enviando una solicitud de carga al
servidor de tabletas.
Esta asignación solo falla si la solicitud de carga de la tableta no se recibe antes de la siguiente
conmutación por error del maestro: un servidor de tabletas solo acepta solicitudes de carga de la
tableta del maestro actual. Por lo tanto, una vez que un maestro envía una solicitud de carga de
tableta, puede asumir que la tableta está asignada hasta que el servidor de tabletas deje de funcionar
o este le informe al maestro que la ha descargado.
Bigtable usa Chubby para monitorear los servidores de tabletas. Cuando un servidor de tabletas
se inicia, crea y adquiere un bloqueo exclusivo en un archivo con un nombre único en un directorio
específico de Chubby. El maestro monitorea este directorio (el directorio de servidores) para descubrir
servidores de tabletas. Un servidor de tabletas deja de servir a sus tabletas si pierde su bloqueo
exclusivo: por ejemplo, una partición de red podría provocar que el servidor pierda su sesión de
Chubby. (Chubby proporciona un mecanismo eficiente que permite a un servidor de tabletas comprobar
si aún mantiene su bloqueo sin generar tráfico de red). Un servidor de tabletas intenta volver a adquirir
un bloqueo exclusivo en su archivo mientras este exista. Si su archivo ya no existe, el servidor de
tabletas no podrá volver a servir, por lo que se autodestruye. Cuando un servidor de tabletas finaliza
(por ejemplo, porque el sistema de administración del clúster se restablece...

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:10 ∙ F. Chang y otros.

al mover la máquina del servidor de tabletas desde el clúster), intenta liberar su bloqueo para que el
maestro reasigne sus tabletas más rápidamente.
El maestro es responsable de detectar cuándo un servidor de tabletas ya no presta servicio a sus
tabletas y de reasignar esas tabletas lo antes posible.
Para detectar cuándo un servidor de tabletas ya no presta servicio a sus tabletas, el maestro solicita
periódicamente a cada servidor el estado de su bloqueo. Si un servidor de tabletas informa que ha
perdido su bloqueo, o si el maestro no ha podido acceder a un servidor en sus últimos intentos, el
maestro intenta adquirir un bloqueo exclusivo en el archivo del servidor. Si el maestro consigue adquirir
el bloqueo, significa que Chubby está activo y el servidor de tabletas está inactivo o tiene problemas
para acceder a él. Por lo tanto, el maestro se asegura de que el servidor de tabletas nunca vuelva a
prestar servicio eliminando su archivo de servidor. Una vez eliminado el archivo de un servidor, el
maestro puede mover todas las tabletas previamente asignadas a ese servidor al conjunto de tabletas
sin asignar. Para garantizar que un clúster de Bigtable no sea vulnerable a problemas de red entre el
maestro y Chubby, el maestro se autodetiene si su sesión de Chubby expira. Los fallos del maestro
no modifican la asignación de tabletas a los servidores de tabletas.

Cuando el sistema de gestión del clúster inicia un maestro, este necesita descubrir las asignaciones
actuales de las tabletas antes de poder modificarlas. El maestro ejecuta los siguientes pasos al iniciarse:

(1) El maestro agarra un candado maestro único en Chubby, que impide la concurrencia
instancias del maestro de alquiler.

(2) El maestro escanea el directorio de servidores en Chubby para encontrar los servidores activos.
(3) El maestro se comunica con cada servidor de tabletas activo para descubrir qué tabletas ya están
asignadas a cada servidor y para actualizar su noción del maestro actual (de modo que cualquier
solicitud de carga de tabletas recibida posteriormente de maestros anteriores será rechazada).

(4) El maestro escanea la tabla METADATOS para conocer el conjunto de tabletas. Al encontrar una
tableta sin asignar, el maestro la añade al conjunto de tabletas sin asignar, lo que la hace elegible
para la asignación.

Una complicación es que el escaneo de la tabla METADATA no puede realizarse hasta que se
hayan asignado las tabletas METADATA . Por lo tanto, antes de iniciar este escaneo (Paso (4)), el
maestro agrega la tableta raíz al conjunto de tabletas sin asignar si no se detectó una asignación para
la tableta raíz durante el Paso (3). Esta adición garantiza que la tableta raíz se asigne. Dado que la
tableta raíz contiene los nombres de todas las tabletas METADATA , el maestro las conoce todas
después de escanearla.

El conjunto de tabletas existentes solo cambia cuando se crea o elimina una tabla, cuando dos
tabletas existentes se fusionan para formar una tableta más grande o cuando una tableta existente se
divide en dos tabletas más pequeñas. El maestro puede realizar un seguimiento de estos cambios
porque inicia todos menos el último. Las divisiones de tabletas se tratan de forma especial, ya que las
inician los servidores de tabletas. Un servidor de tabletas confirma una división registrando la
información de la nueva tableta en la tabla METADATA . Tras confirmar la división, el servidor de
tabletas notifica al maestro. Si la notificación de división es...
ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido para datos estructurados ∙ 4:11

Fig. 6. Representación en tableta.

Si se pierde (debido a la falla del servidor de tabletas o del maestro), el maestro detecta la nueva tableta
al solicitar a un servidor de tabletas que cargue la tableta que se ha dividido. El servidor de tabletas
notificará al maestro sobre la división, ya que la entrada de tableta que encuentra en la tabla METADATA
solo especificará una parte de la tableta que el maestro le solicitó cargar.

5.3 Servicio de tabletas

El estado persistente de una tableta se almacena en GFS, como se ilustra en la Figura 6.


Las actualizaciones se registran en un registro de confirmaciones que almacena los registros de rehacer.
Las confirmaciones recientes se almacenan en memoria, en un búfer ordenado llamado tabla de memoria.
Una tabla de memoria mantiene las actualizaciones fila por fila, donde cada fila se copia al escribir para
mantener la consistencia a nivel de fila. Las actualizaciones anteriores se almacenan en una secuencia
de tablas SSTable (inmutables).
Para recuperar una tableta, un servidor lee sus metadatos de la tabla METADATA . Estos metadatos
contienen la lista de SSTables que componen la tableta y un conjunto de puntos de rehacer, que apuntan
a cualquier registro de confirmación que pueda contener datos de la tableta. El servidor lee los índices de
las SSTables en memoria y reconstruye la memtable aplicando todas las actualizaciones confirmadas
desde los puntos de rehacer.

Cuando una operación de escritura llega a un servidor de tabletas, este comprueba que esté
correctamente formada (es decir, que no provenga de un cliente con errores u obsoleto) y que el remitente
esté autorizado para realizar la mutación. La autorización se realiza leyendo la lista de escritores permitidos
de un archivo Chubby (que casi siempre coincide con la caché del cliente Chubby). Una mutación válida
se escribe en el registro de confirmaciones.
La confirmación grupal se utiliza para mejorar el rendimiento de pequeñas mutaciones [DeWitt et al., 1984;
Gawlick y Kinkade, 1985]. Una vez confirmada la escritura, su contenido se inserta en la tabla de memoria.

Cuando una operación de lectura llega a un servidor de tabletas, se verifica de forma similar su
correcta formación y autorización. Una operación de lectura válida se ejecuta en una vista fusionada de la
secuencia de SSTables y la tabla de memoria. Dado que
ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:12
∙ F. Chang y otros.

SSTables y memtable son estructuras de datos ordenadas lexicográficamente, la vista fusionada


se puede formar de manera eficiente.
Las operaciones de lectura y escritura entrantes pueden continuar mientras las tabletas se
dividen y fusionan. También pueden ocurrir operaciones de lectura y escritura mientras se
compactan las tabletas; las compactaciones se describen en la siguiente sección.

5.4 Compactaciones.
A medida que se ejecutan las operaciones de escritura, el tamaño de la tabla de memoria aumenta.
Cuando alcanza un umbral, se congela, se crea una nueva tabla y la tabla congelada se convierte
en una SSTable y se escribe en GFS. Este pequeño proceso de compactación tiene dos objetivos:
reduce el uso de memoria del servidor de la tableta y la cantidad de datos que deben leerse del
registro de confirmaciones durante la recuperación si el servidor falla.

Cada compactación menor crea una nueva SSTable. Si este comportamiento persiste sin
control, las operaciones de lectura podrían necesitar fusionar actualizaciones de un número
arbitrario de SSTables. En su lugar, limitamos el número de estos archivos ejecutando
periódicamente una compactación de fusión en segundo plano. Una compactación de fusión lee
el contenido de algunas SSTables y la memtable, y escribe una nueva SSTable. Las SSTables y
la memtable de entrada se pueden descartar una vez finalizada la compactación.

Una compactación por fusión que reescribe todas las SSTables en una sola se denomina
compactación mayor. Las SSTables generadas por compactaciones no mayores pueden contener
entradas de eliminación especiales que suprimen los datos eliminados en las SSTables antiguas
que aún están activas. Por otro lado, una compactación mayor produce una SSTable que no
contiene información ni datos eliminados. Bigtable recorre todas sus tabletas y les aplica
compactaciones mayores regularmente. Estas compactaciones mayores permiten a Bigtable
recuperar los recursos utilizados por los datos eliminados y garantizar que estos desaparezcan
del sistema de forma oportuna, lo cual es importante para los servicios que almacenan datos
confidenciales.
El rendimiento de lectura de Bigtable se beneficia de una optimización de localidad en GFS.
Al escribir archivos, GFS intenta colocar una réplica de los datos en la misma máquina que el
escritor. Al leer archivos GFS, las lecturas se realizan desde la réplica más cercana disponible.
Por lo tanto, en el caso habitual de servidores de tabletas que comparten máquinas con servidores
GFS, estos compactan los datos en SSTables que tienen una réplica en el disco local, lo que
permite un acceso rápido a dichas SSTables al procesar solicitudes de lectura posteriores.

5.5 Gestión de esquemas. Los


esquemas de Bigtable se almacenan en Chubby. Chubby es un sustrato de comunicación eficaz
para los esquemas de Bigtable, ya que proporciona escrituras atómicas de archivos completos y
almacenamiento en caché consistente de archivos pequeños. Por ejemplo, supongamos que un
cliente desea eliminar algunas familias de columnas de una tabla. El maestro realiza
comprobaciones de control de acceso, verifica que el esquema resultante esté correctamente
formado y, a continuación, instala el nuevo esquema reescribiendo el archivo de esquema
correspondiente en Chubby. Cuando los servidores de tabletas necesitan determinar qué familias
de columnas existen, simplemente leen el archivo de esquema correspondiente de Chubby, que casi siempre está disponible en
ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido para datos estructurados ∙ 4:13

La caché del cliente Chubby del servidor. Dado que las cachés de Chubby son consistentes, se garantiza
que los servidores de tabletas verán todos los cambios en ese archivo.

6. MEJORAS

La implementación descrita en la sección anterior requirió varias mejoras para lograr el rendimiento, la
disponibilidad y la confiabilidad requeridos por nuestros usuarios. Esta sección describe partes de la
implementación con más detalle para destacar estas mejoras.

Grupos de localidad. Cada familia de columnas se asigna a un grupo de localidad definido por el
cliente, que es una abstracción que permite a los clientes controlar la distribución de almacenamiento de
sus datos. Durante la compactación, se genera una SSTable independiente para cada grupo de localidad
en cada tableta. Separar las familias de columnas a las que no se suele acceder juntas en grupos de
localidad independientes permite lecturas más eficientes.
Por ejemplo, los metadatos de la página en Webtable (como el idioma y las sumas de comprobación)
pueden estar en un grupo de localidad y el contenido de la página puede estar en un grupo diferente; una
aplicación que desea leer los metadatos no necesita leer todo el contenido de la página.

Además, se pueden especificar algunos parámetros de ajuste útiles para cada grupo de localidades.
Por ejemplo, se puede declarar que un grupo de localidades está en memoria.
Las SSTables para grupos de localidad en memoria se cargan de forma diferida en la memoria del
servidor de la tableta. Dado que las SSTables son inmutables, no hay problemas de consistencia. Una
vez cargadas, las familias de columnas que pertenecen a dichos grupos de localidad se pueden leer sin
acceder al disco. Esta función es útil para pequeños fragmentos de datos a los que se accede con
frecuencia; la utilizamos internamente para la familia de columnas de ubicación de la tableta en la tabla
METADATA .

Compresión. Los clientes pueden controlar si las SSTables de un grupo de localidades se comprimen
o no y, en tal caso, el formato de compresión utilizado. El formato de compresión especificado por el
usuario se aplica a cada bloque de SSTable (cuyo tamaño se puede controlar mediante un parámetro de
ajuste específico del grupo de localidades). Aunque se pierde espacio en disco al comprimir cada bloque
por separado, la ventaja es que se pueden leer pequeñas porciones de una SSTable sin descomprimir el
archivo completo. Muchos clientes de Bigtable utilizan un esquema de compresión personalizado de dos
pasadas. La primera utiliza el esquema de Bentley y McIlroy [1999], que comprime cadenas largas
comunes en una ventana grande. La segunda utiliza un algoritmo de compresión rápido que busca
repeticiones en una pequeña ventana de datos de 16 KB. Ambas pasadas de compresión son muy
rápidas: codifican a 100­200 MB/s y decodifican a 400­1000 MB/s en equipos modernos.

Aunque priorizamos la velocidad en lugar de la reducción de espacio al elegir nuestros algoritmos de


compresión, este esquema de compresión de dos pasadas funciona sorprendentemente bien. Por
ejemplo, en Webtable, utilizamos este esquema de compresión para almacenar el contenido de páginas
web. En un experimento, almacenamos una gran cantidad de documentos en un grupo de localidad
comprimido. Para los fines del experimento, nos limitamos a una versión de cada documento en lugar de
almacenar todas las versiones disponibles. El esquema logró una reducción de espacio de 10 a 1. Esto
es mucho mejor que las reducciones típicas de Gzip de 3 a 1 o 4 a 1 en páginas HTML porque

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:14
∙ F. Chang y otros.

La disposición de las filas de Webtable se debe a que todas las páginas de un mismo host se
almacenan cerca unas de otras. Esto permite que el algoritmo Bentley­McIlroy identifique grandes
cantidades de código repetitivo compartido en páginas del mismo host. Muchas aplicaciones, no solo
Webtable, eligen los nombres de sus filas de forma que los datos similares se agrupen y, por lo tanto,
logren excelentes tasas de compresión. Las tasas de compresión mejoran aún más al almacenar
varias versiones del mismo valor en Bigtable.

Almacenamiento en caché para mejorar el rendimiento de lectura. Para mejorar el rendimiento de


lectura, los servidores de tabletas utilizan dos niveles de almacenamiento en caché. La caché de
escaneo es una caché de nivel superior que almacena en caché los pares clave­valor devueltos por la
interfaz SSTable al código del servidor de tabletas. La caché de bloque es una caché de nivel inferior
que almacena en caché los bloques SSTable leídos desde GFS. La caché de escaneo es especialmente
útil para aplicaciones que suelen leer los mismos datos repetidamente. La caché de bloque es útil para
aplicaciones que suelen leer datos similares a los leídos recientemente (por ejemplo, lecturas
secuenciales o lecturas aleatorias de diferentes columnas en el mismo grupo de localidad dentro de
una fila activa).

Filtros Bloom. Como se describe en la Sección 5.3, una operación de lectura debe leer de todas
las SSTables que conforman el estado de una tableta. Si estas SSTables no están en memoria,
podríamos tener que realizar muchos accesos al disco. Reducimos el número de accesos permitiendo
a los clientes especificar la creación de filtros Bloom [Bloom 1970] para las SSTables de un grupo de
localidades específico. Un filtro Bloom permite preguntar si una SSTable puede contener datos para
un par de filas/columnas especificado. En ciertas aplicaciones, la pequeña cantidad de memoria del
servidor de tabletas utilizada para almacenar filtros Bloom reduce drásticamente el número de
búsquedas en disco necesarias para las operaciones de lectura. El uso de filtros Bloom también evita
los accesos al disco en la mayoría de las búsquedas de filas o columnas inexistentes.

Implementación del registro de confirmación. Si se mantuviera el registro de confirmación de cada


tableta en un archivo de registro independiente, se escribiría una gran cantidad de archivos
simultáneamente en GFS. Dependiendo de la implementación del sistema de archivos subyacente en
cada servidor GFS, estas escrituras podrían requerir un gran número de búsquedas en disco para
escribir en los diferentes archivos de registro físicos. Además, tener archivos de registro separados
por tableta también reduce la eficacia de la optimización de confirmación grupal, ya que los grupos
tenderían a ser más pequeños. Para solucionar estos problemas, se añaden mutaciones a un único
registro de confirmación por servidor de tableta, lo que combina las mutaciones de diferentes tabletas
en el mismo archivo de registro físico [Gray 1978; Hagmann 1987].
Usar un solo registro ofrece importantes mejoras de rendimiento durante el funcionamiento normal,
pero dificulta la recuperación. Cuando un servidor de tabletas deja de funcionar, las tabletas que servía
se migrarán a un gran número de otros servidores; cada servidor suele cargar una pequeña cantidad
de tabletas del servidor original. Para recuperar el estado de una tableta, el nuevo servidor debe volver
a aplicar las mutaciones de esa tableta desde el registro de confirmación escrito por el servidor original.
Sin embargo, las mutaciones de estas tabletas se combinaron en el mismo archivo de registro físico.

Un enfoque sería que cada nuevo servidor de tabletas lea este archivo de registro de confirmación
completo y aplique solo las entradas necesarias para las tabletas que necesita recuperar. Sin
embargo, bajo este esquema, si a 100 máquinas se les asignara una sola tableta cada una...
ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido para datos estructurados ∙ 4:15

desde un servidor de tableta fallido, entonces el archivo de registro se leería 100 veces (una vez por
cada servidor).
Evitamos la duplicación de lecturas de registros ordenando primero las entradas del registro de
confirmación por tabla de claves, nombre de fila y número de secuencia del registro. En la salida
ordenada, todas las mutaciones de una tableta en particular son contiguas y, por lo tanto, se pueden
leer eficientemente con una búsqueda en disco seguida de una lectura secuencial. Para paralelizar la
ordenación, particionado el archivo de registro en segmentos de 64 MB, ordenando cada segmento
en paralelo en diferentes servidores de tabletas. Este proceso de ordenación es coordinado por el
servidor maestro y se inicia cuando un servidor de tabletas indica que necesita recuperar mutaciones
de algún archivo de registro de confirmación.
La escritura de registros de confirmación en GFS a veces causa interrupciones en el rendimiento
por diversas razones (por ejemplo, un servidor GFS involucrado en la escritura se bloquea, tiene
mucha carga o las rutas de red utilizadas para acceder al conjunto específico de servidores GFS
sufren congestión). Para proteger las mutaciones de los picos de latencia de GFS, cada servidor de
tableta cuenta con dos subprocesos de escritura de registros, cada uno de los cuales escribe en su
propio archivo de registro; solo uno de estos subprocesos está activo a la vez. Si las escrituras en el
archivo de registro activo tienen un rendimiento deficiente, la escritura del archivo de registro se
transfiere al otro subproceso, y las mutaciones en la cola del registro de confirmación son escritas por
el nuevo subproceso de escritura de registros activo. Las entradas de registro contienen números de
secuencia para que el proceso de recuperación pueda eliminar las entradas duplicadas resultantes de
este proceso de intercambio de registros.

Aceleración de la recuperación de la tableta. Antes de descargar una tableta, el servidor realiza


una compactación menor. Esta compactación reduce el tiempo de recuperación al reducir la cantidad
de estado sin compactar en el registro de confirmación del servidor. Tras finalizar esta compactación,
el servidor deja de servirla. Antes de descargarla, realiza otra compactación menor (normalmente muy
rápida) para eliminar cualquier estado sin compactar restante en el registro del servidor que se haya
generado durante la primera compactación menor. Tras esta segunda compactación menor, la tableta
puede cargarse en otro servidor sin necesidad de recuperar las entradas del registro.

Aprovechamiento de la inmutabilidad. Además de las cachés de SSTable, se han simplificado otras


partes del sistema Bigtable gracias a que todas las SSTables que generamos son inmutables. Por
ejemplo, no necesitamos sincronizar los accesos al sistema de archivos al leer desde las SSTables.
Como resultado, el control de concurrencia sobre las filas se puede implementar de forma muy
eficiente. La única estructura de datos mutable a la que se accede tanto en lecturas como en escrituras
es la memtable. Para reducir la contención durante las lecturas de la memtable, configuramos cada
fila de la memtable como copia al escribir y permitimos que las lecturas y escrituras se realicen en
paralelo.
Dado que las SSTables son inmutables, el problema de eliminar de forma permanente los datos
eliminados se transforma en el de recolectar basura de las SSTables obsoletas.
Las SSTables de cada tableta se registran en la tabla METADATA . El maestro elimina las SSTables
obsoletas mediante una recolección de basura de marcado y barrido [McCarthy 1960] del conjunto de
SSTables, donde la tabla METADATA contiene el conjunto de raíces.

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:16
∙ F. Chang y otros.

Finalmente, la inmutabilidad de SSTables nos permite dividir tabletas rápidamente.


En lugar de generar un nuevo conjunto de SSTables para cada tableta secundaria, permitimos que las
tabletas secundarias compartan las SSTables de la tableta principal.

7. EVALUACIÓN DEL DESEMPEÑO

Configuramos un clúster de Bigtable con N servidores tablet para medir el rendimiento y la escalabilidad
de Bigtable a medida que N varía. Los servidores tablet se configuraron para usar 1 GB de memoria y
escribir en una celda GFS compuesta por 1786 máquinas con dos discos duros IDE de 400 GB cada
una. N máquinas cliente generaron la carga de Bigtable utilizada para estas pruebas. (Usamos el mismo
número de clientes que servidores tablet para garantizar que los clientes nunca fueran un cuello de
botella). Cada máquina tenía dos chips Opteron de doble núcleo a 2 GHz, memoria física suficiente
para albergar el conjunto de trabajo de todos los procesos en ejecución y un único enlace Gigabit
Ethernet. Las máquinas se organizaron en una red conmutada de dos niveles en forma de árbol con
aproximadamente 100­200 Gb/s de ancho de banda agregado disponible en la raíz. Todas las máquinas
se encontraban en la misma instalación de alojamiento y, por lo tanto, el tiempo de ida y vuelta entre
cualquier par de máquinas era inferior a un milisegundo.

Los servidores de tabletas, el maestro, los clientes de prueba y los servidores GFS se ejecutaban en
el mismo conjunto de máquinas. Cada máquina ejecutaba un servidor GFS. Algunas máquinas también
ejecutaban un servidor de tabletas, un proceso cliente o procesos de otros trabajos que usaban el grupo
simultáneamente con estos experimentos.
R es el número específico de claves de fila de Bigtable involucradas en la prueba. Se eligió R para
que cada punto de referencia leyera o escribiera aproximadamente 1 GB de datos por servidor de tableta.

El benchmark de escritura secuencial utilizó claves de fila con nombres de 0 a R − 1. Este espacio
de claves de fila se dividió en 10N rangos de igual tamaño. Estos rangos fueron asignados a los N
clientes por un programador central que asignaba el siguiente rango disponible a un cliente tan pronto
como este terminaba de procesar el rango anterior. Esta asignación dinámica ayudó a mitigar los
efectos de las variaciones de rendimiento causadas por otros procesos en ejecución en los equipos
cliente. Se escribió una sola cadena bajo cada clave de fila. Cada cadena se generó aleatoriamente y,
por lo tanto, no era comprimible. Además, las cadenas bajo diferentes claves de fila eran distintas, por
lo que no fue posible la compresión entre filas.

La prueba de referencia de escritura aleatoria fue similar, excepto que la clave de fila se sometió a un
algoritmo hash módulo R inmediatamente antes de escribir, de modo que la carga de escritura se
distribuyó de manera casi uniforme en todo el espacio de fila durante toda la duración de la prueba de referencia.
La prueba de lectura secuencial generó claves de fila exactamente igual que la prueba de escritura
secuencial, pero en lugar de escribir bajo la clave de fila, leyó la cadena almacenada bajo ella (escrita
mediante una invocación previa de la prueba de escritura secuencial). De igual forma, la prueba de
lectura aleatoria replicó la operación de la prueba de escritura aleatoria.

El análisis comparativo es similar al análisis comparativo de lectura secuencial, pero utiliza la


compatibilidad de la API de Bigtable para analizar todos los valores de un rango de filas. El análisis
reduce el número de RPC ejecutados por el análisis comparativo, ya que un solo RPC obtiene una gran
secuencia de valores de una tableta.
servidor.

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido para datos estructurados ∙ 4:17

Tabla I. Número de valores de 1000 bytes leídos/escritos por


segundo. Los valores corresponden a la tarifa por servidor tablet.

Experimento Número de servidores de tabletas


1 50 250 500
lecturas aleatorias 1212 593 479 241
lecturas aleatorias (mem) 10811 8511 8000 6250
escrituras aleatorias 8850 3745 3425 2000
lecturas secuenciales 4425 2463 2625 2469
escrituras secuenciales 8547 3623 2451 1905
escaneos 15385 10526 9524 7843

Fig. 7. Número de valores de 1000 bytes leídos/escritos por segundo. Las curvas indican el total.
tasa en todos los servidores de tabletas.

El punto de referencia de lecturas aleatorias (mem) es similar al punto de referencia de lecturas


aleatorias, pero el grupo de localidades que contiene los datos de referencia está marcado como
en memoria, por lo que las lecturas se realizan desde la memoria del servidor de la tableta en lugar
de requerir una lectura GFS. Solo para esta prueba de rendimiento, redujimos la
cantidad de datos por servidor de tableta de 1 GB a 100 MB para que quepa
cómodamente en la memoria disponible en el servidor de la tableta.
La Tabla I y la Figura 7 ofrecen dos perspectivas sobre el rendimiento de nuestros puntos de
referencia al leer y escribir valores de 1000 bytes en Bigtable. La tabla...
muestra el número de operaciones por segundo por servidor de tableta; el gráfico muestra
el número agregado de operaciones por segundo.

Rendimiento de una sola tableta­servidor. Analicemos primero el rendimiento con


Solo un servidor de tableta. Las lecturas aleatorias son más lentas que todas las demás operaciones.
orden de magnitud o más. Cada lectura aleatoria implica la transferencia de 64 KB.
Bloque SSTable a través de la red desde GFS a un servidor de tableta, de los cuales solo
Se utiliza un único valor de 1000 bytes. El servidor de la tableta ejecuta aproximadamente
1200 lecturas por segundo, lo que se traduce en aproximadamente 75 MB/s de datos
Leer desde GFS. Este ancho de banda es suficiente para saturar las CPU del servidor de la tableta.
Debido a los costos generales en nuestra pila de red, análisis de SSTable y Bigtable
código, y también es casi suficiente para saturar los enlaces de red utilizados en nuestro sistema.
La mayoría de las aplicaciones Bigtable con este tipo de patrón de acceso reducen la
tamaño de bloque a un valor más pequeño, normalmente 8 KB.

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:18
∙ F. Chang y otros.

Las lecturas aleatorias de la memoria son mucho más rápidas ya que cada lectura de 1000 bytes es
satisfecho con la memoria local del servidor de la tableta sin tener que recuperar un gran volumen de 64 KB
Bloque de GFS.
Las lecturas secuenciales funcionan mejor que las lecturas aleatorias, ya que cada 64 KB
El bloque SSTable que se obtiene de GFS se almacena en nuestro caché de bloques, donde
Se utiliza para atender las siguientes 64 solicitudes de lectura.
Los escaneos son aún más rápidos ya que el servidor de la tableta puede devolver una gran cantidad
de valores en respuesta a una sola RPC de cliente y, por lo tanto, la sobrecarga de RPC es
amortizado en un gran número de valores.
Las escrituras funcionan mejor que las lecturas porque cada servidor de tableta agrega todos los
Las escrituras entrantes se realizan en un único registro de confirmación y, debido a que usamos la confirmación grupal para
Transmite estas escrituras eficientemente a GFS. Las lecturas, por otro lado, suelen requerir una búsqueda de
disco por cada SSTable accedida. Las escrituras aleatorias y secuenciales tienen un rendimiento muy similar; en
ambos casos, todas las escrituras a...
Los servidores de tabletas se registran en el mismo registro de confirmaciones.

Escalado. El rendimiento agregado aumenta drásticamente, en más de un factor de


cien, a medida que aumentamos el número de servidores de tabletas en el sistema de 1
a 500. Por ejemplo, el rendimiento de las lecturas aleatorias de la memoria aumenta
en casi un factor de 300 a medida que el número de servidores de tabletas aumenta en un factor
de 500. Este comportamiento se produce porque el cuello de botella en el rendimiento de este
El punto de referencia es la CPU del servidor de tableta individual.
Sin embargo, el rendimiento no aumenta linealmente. Para la mayoría de los puntos de referencia,
Hay una caída significativa en el rendimiento por servidor cuando se pasa de 1 a 50.
Servidores de tabletas. Esta caída se debe a un desequilibrio en la carga en varias configuraciones de servidores,
a menudo debido a que otros procesos compiten por la CPU y la red.
El algoritmo de equilibrio de carga intenta solucionar este desequilibrio, pero no puede hacerlo.
Un trabajo perfecto por dos razones principales: el reequilibrio se limita para reducir el número de movimientos de
la tableta (una tableta no está disponible durante un corto período de tiempo, generalmente menos
más de un segundo, cuando se mueve), y la carga generada por nuestros puntos de referencia
cambia a medida que avanza el índice de referencia.
El benchmark de lectura aleatoria muestra el peor escalamiento (un aumento del rendimiento agregado de
solo un factor de 100 para un aumento de 500 veces en el número de servidores). Este comportamiento se debe
a que (como se explicó anteriormente) transferimos
Un bloque grande de 64 KB por cada 1000 bytes leídos en la red. Esta transferencia satura varios enlaces
compartidos de 1 Gigabit en nuestra red y, como resultado,
El rendimiento por servidor disminuye significativamente a medida que aumentamos el número de
máquinas.

8. APLICACIONES REALES

A partir de agosto de 2006, hay 388 clústeres de Bigtable que no son de prueba ejecutándose en varios
Clústeres de máquinas de Google, con un total combinado de aproximadamente 24.500 servidores de tabletas.
La Tabla II muestra una distribución aproximada de los servidores de tabletas por clúster. Muchos de estos...
Los clústeres se utilizan para fines de desarrollo y, por lo tanto, permanecen inactivos durante períodos
considerables. Un grupo de 14 clústeres activos con un total de 8069 servidores de tabletas.
registró un volumen agregado de más de 1,2 millones de solicitudes por segundo, con

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido para datos estructurados ∙ 4:19

Tabla II. Distribución del número de servidores de tabletas en clústeres de Bigtable.


# de servidores de tabletas # de clústeres
.. 19 259
0 20 .. 49 47
50 .. 99 20
100 .. > 499 50
500 12

Tabla III. Características de algunas mesas en uso productivo.

Proyecto Tamaño Comp. # (B) # # % Interfaz


nombre Relación (TB) Células Familias Grupos MMap
Gatear 800 11% 1000 16 8 0% No
Gatear 50 33% 200 2 2 0% No
Analítica 20 29% 10 1 1 0% Sí
Analítica 200 14% 80 0% Sí
Base 2 31% 10 1 29 13 15% Sí
Tierra 0.5 64% 70 8 7 2 33% Sí
Tierra – 9 8 3 0% No
Orkut 9 – 0.9 8 5 1% Sí
Búsqueda personal 4 47% 6 93 11 5% Sí

El tamaño (medido antes de la compresión) y el número de celdas indican tamaños aproximados. La relación de
compresión no se proporciona para las tablas con la compresión deshabilitada. El frontend indica que
El rendimiento de la aplicación es sensible a la latencia.

tráfico RPC entrante de aproximadamente 741 MB/s y tráfico RPC saliente de aproximadamente
16 GB/s.
La Tabla III proporciona algunos datos sobre algunas de las tablas en uso en agosto.
2006. Algunas tablas almacenan datos que se proporcionan a los usuarios, mientras que otras almacenan datos
para el procesamiento por lotes; las tablas varían ampliamente en tamaño total, tamaño de celda promedio,
porcentaje de datos servidos desde la memoria y complejidad del esquema de la tabla.
En el resto de esta sección, describimos brevemente cómo tres equipos de productos utilizan
Mesa grande.

8.1 Google Analytics

Google Analytics ([Link]) es un servicio que ayuda a los webmasters a analizar los patrones de
tráfico en sus sitios web. Proporciona estadísticas agregadas, como
como el número de visitantes únicos por día y las páginas vistas por URL por día,
así como informes de seguimiento del sitio, como el porcentaje de usuarios que realizaron una
compra, dado que anteriormente vieron una página específica.
Para habilitar el servicio, los webmasters incorporan un pequeño programa JavaScript en
sus páginas web. Este programa se invoca cada vez que se visita una página. Registra
Información diversa sobre la solicitud en Google Analytics, como el identificador del usuario e información sobre
la página que se está recuperando. Google Analytics resume estos datos y los pone a disposición de los
webmasters.
Describiremos brevemente dos de las tablas que utiliza Google Analytics. El clic sin procesar
La tabla (~200 TB) mantiene una fila para cada sesión de usuario final. El nombre de la fila es
una tupla que contiene el nombre del sitio web y la hora en que se realizó la sesión
creado. Este esquema garantiza que las sesiones que visitan el mismo sitio web sean

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:20 ∙ F. Chang y otros.

Contiguos y ordenados cronológicamente. Esta tabla se comprime al 14 % de su tamaño original.

La tabla de resumen (~20 TB) contiene varios resúmenes predefinidos para cada sitio web. Esta tabla
se genera a partir de la tabla de clics sin procesar mediante trabajos de MapReduce programados
periódicamente. Cada trabajo de MapReduce extrae datos de sesiones recientes de la tabla de clics sin
procesar. El rendimiento general del sistema está limitado por el rendimiento de GFS. Esta tabla se
comprime al 29 % de su tamaño original.

8.2 Google Earth.

Google opera una colección de servicios que proporcionan a los usuarios imágenes satelitales de alta
resolución de la superficie terrestre. Los usuarios pueden acceder a las imágenes a través de la interfaz
web de Google Maps ([Link]) y del software cliente personalizado de Google Earth
([Link]) . Estos productos permiten a los usuarios navegar por la superficie terrestre: pueden
desplazarse, visualizar y anotar imágenes satelitales con distintos niveles de resolución. Este sistema
utiliza una tabla para preprocesar los datos y un conjunto diferente de tablas para procesar los datos del
cliente.
El proceso de preprocesamiento utiliza una tabla para almacenar imágenes sin procesar. Durante el
preprocesamiento, las imágenes se depuran y se consolidan para obtener los datos finales de la publicación.
Esta tabla contiene aproximadamente 70 terabytes de datos y, por lo tanto, se sirve desde el disco. Las
imágenes ya están comprimidas eficientemente, por lo que la compresión de Bigtable está deshabilitada.

Cada fila de la tabla de imágenes corresponde a un solo segmento geográfico.


Las filas se nombran para garantizar que los segmentos geográficos adyacentes se almacenen cerca. La
tabla contiene una familia de columnas para registrar las fuentes de datos de cada segmento. Esta familia
de columnas tiene un gran número de columnas: básicamente, una por cada imagen de datos sin procesar.
Dado que cada segmento se construye a partir de unas pocas imágenes, esta familia de columnas es muy
dispersa.
El flujo de preprocesamiento depende en gran medida de MapReduce sobre Bigtable para transformar
los datos. El sistema procesa más de 1 MB/s de datos por servidor de tableta durante algunas de estas
tareas de MapReduce.
El sistema de servicio utiliza una tabla para indexar los datos almacenados en GFS. Esta tabla es
relativamente pequeña (aproximadamente 500 GB), pero debe procesar decenas de miles de consultas
por segundo por centro de datos con baja latencia. Por lo tanto, esta tabla se aloja en cientos de servidores
de tabletas y contiene familias de columnas en memoria.

8.3 Búsqueda personalizada

Búsqueda personalizada ([Link]/psearch) es un servicio opcional que registra las consultas y


los clics de los usuarios en diversas propiedades de Google, como búsquedas web, imágenes y noticias.
Los usuarios pueden explorar su historial de búsqueda para revisar sus consultas y clics anteriores, y
solicitar resultados de búsqueda personalizados según sus patrones de uso de Google.

La Búsqueda Personalizada almacena los datos de cada usuario en Bigtable. Cada usuario tiene un ID
de usuario único y se le asigna una fila con ese ID. Todas las acciones del usuario se almacenan en una
tabla. Cada tipo de acción tiene una familia de columnas independiente (por ejemplo, hay una familia de
columnas que almacena todas las consultas web). Cada elemento de datos utiliza como marca de tiempo
de Bigtable la hora en que se realizó la acción del usuario correspondiente.
ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido para datos estructurados ∙ 4:21

Ocurrió. La Búsqueda Personalizada genera perfiles de usuario mediante MapReduce sobre


Bigtable. Estos perfiles se utilizan para personalizar los resultados de búsqueda en tiempo real.
Los datos de Búsqueda Personalizada se replican en varios clústeres de Bigtable para
aumentar la disponibilidad y reducir la latencia debido a la distancia de los clientes. El equipo de
Búsqueda Personalizada desarrolló inicialmente un mecanismo de replicación del lado del cliente
sobre Bigtable que garantizaba la consistencia final de todas las réplicas. El sistema actual utiliza
un subsistema de replicación integrado en los servidores.
El diseño del sistema de almacenamiento de Búsqueda Personalizada permite que otros
grupos agreguen información por usuario en sus propias columnas. Actualmente, el sistema lo
utilizan muchas otras propiedades de Google que necesitan almacenar opciones y ajustes de
configuración por usuario. Compartir una tabla entre varios grupos resultó en un número
inusualmente grande de familias de columnas. Para facilitar el uso compartido, añadimos un
mecanismo de cuotas simple a Bigtable para limitar el consumo de almacenamiento de cada
cliente en las tablas compartidas. Este mecanismo proporciona cierto aislamiento entre los
distintos grupos de productos que utilizan este sistema para el almacenamiento de información
por usuario.

9. LECCIONES

En el proceso de diseño, implementación, mantenimiento y soporte de Bigtable, adquirimos


experiencia útil y aprendimos varias lecciones interesantes.
Una lección que aprendimos es que los grandes sistemas distribuidos son vulnerables a
muchos tipos de fallos, no solo a las particiones de red estándar y los fallos de parada por error
que se asumen en muchos protocolos distribuidos. Por ejemplo, hemos visto problemas debido a
las siguientes causas: corrupción de memoria y red, gran desfase de reloj, máquinas colgadas,
particiones de red extendidas y asimétricas, errores en otros sistemas que utilizamos (por
ejemplo, Chubby), desbordamiento de las cuotas de GFS y mantenimiento de hardware planificado
y no planificado. A medida que adquirimos más experiencia con estos problemas, los hemos
abordado modificando varios protocolos. Por ejemplo, añadimos suma de comprobación a nuestro
mecanismo de RPC. También gestionamos algunos problemas eliminando las suposiciones que
una parte del sistema hacía sobre otra. Por ejemplo, dejamos de asumir que una operación
determinada de Chubby podía devolver solo uno de un conjunto fijo de errores.

Otra lección que aprendimos es que es importante retrasar la incorporación de nuevas


funciones hasta que esté claro cómo se utilizarán. Por ejemplo, inicialmente planeamos admitir
transacciones de propósito general en nuestra API. Sin embargo, dado que no teníamos un uso
inmediato para ellas, no las implementamos.
Ahora que contamos con muchas aplicaciones reales ejecutándose en Bigtable, hemos podido
examinar sus necesidades reales y hemos descubierto que la mayoría solo requieren transacciones
de una sola fila. En los casos en que se han solicitado transacciones distribuidas, el uso más
importante es el mantenimiento de índices secundarios, y planeamos añadir un mecanismo
especializado para satisfacer esta necesidad. El nuevo mecanismo será menos general que las
transacciones distribuidas, pero será más eficiente (especialmente para actualizaciones que
abarcan cientos de filas o más) y también interactuará mejor con nuestro esquema de replicación
optimista entre centros de datos.
Una lección práctica que aprendimos al dar soporte a Bigtable es la importancia de una
monitorización adecuada a nivel de sistema (es decir, la monitorización tanto de Bigtable como de
los procesos del cliente que lo utilizan). Por ejemplo, ampliamos nuestra
ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:22 ∙ F. Chang y otros.

Sistema RPC para mantener seguimientos detallados de una muestra de RPC. Esta función tiene
Nos permitió detectar y solucionar muchos problemas como la contención de bloqueo en la tableta.
estructuras de datos, escrituras lentas en GFS al confirmar mutaciones de Bigtable,
y accesos bloqueados a la tabla METADATA cuando las tabletas METADATA no están disponibles.
Otro ejemplo de monitoreo útil es que cada clúster de Bigtable está...
registrados en Chubby. Esto nos permite rastrear todos los clústeres y descubrir cómo...
grandes que son, ver qué versiones de nuestro software están ejecutando, cuánto
tráfico que están recibiendo y si hay o no problemas como
Latencias inesperadamente grandes.
La lección más importante que aprendimos es el valor de los diseños simples. Dado
tanto el tamaño de nuestro sistema (aproximadamente 100.000 líneas de código que no son de prueba), como la
El hecho de que el código evolucione con el tiempo de maneras inesperadas nos ha llevado a descubrir que el código...
Y la claridad del diseño son de gran ayuda en el mantenimiento y la depuración del código.
Un ejemplo de esto es nuestro protocolo de membresía para servidores de tabletas. Nuestro primer protocolo
era simple: el maestro emitía periódicamente concesiones a los servidores de tabletas, y
Los servidores de tabletas se desactivaban automáticamente si su contrato de arrendamiento expiraba.
Desafortunadamente, este protocolo reducía significativamente la disponibilidad ante problemas de red.
y también era sensible al tiempo de recuperación del maestro. Rediseñamos el protocolo.
Varias veces hasta que obtuvimos un protocolo que funcionó bien. Sin embargo, el protocolo resultante era
demasiado complejo y dependía del comportamiento de las funciones de Chubby, que rara vez utilizaban otras
aplicaciones. Descubrimos que...
Estaban gastando una cantidad excesiva de tiempo depurando casos especiales oscuros,
no solo en el código de Bigtable, sino también en el de Chubby. Finalmente, descartamos
Este protocolo y se trasladó a un protocolo más nuevo y más simple que depende únicamente de
Funciones de Chubby ampliamente utilizadas.

10. TRABAJOS RELACIONADOS

El proyecto Boxwood [MacCormick et al. 2004] tiene componentes que se superponen en algunos aspectos con
Chubby, GFS y Bigtable, ya que proporciona
acuerdo distribuido, bloqueo, almacenamiento de fragmentos distribuidos y distribuido
Almacenamiento en árbol B. En cada caso donde hay solapamiento, parece que el componente de Boxwood
está orientado a un nivel ligeramente inferior al del servicio de Google correspondiente. El objetivo del proyecto
Boxwood es proporcionar infraestructura.
para construir servicios de nivel superior, como sistemas de archivos o bases de datos, mientras que
El objetivo de Bigtable es dar soporte directo a las aplicaciones cliente que deseen almacenar
datos.
Muchos proyectos recientes han abordado el problema de proporcionar almacenamiento distribuido o
servicios de nivel superior a través de redes de área amplia, a menudo a “escala de Internet”.
Esto incluye el trabajo en tablas hash distribuidas que comenzó con proyectos como
como CAN [Ratnasamy et al. 2001], Acorde [Stoica et al. 2001], Tapiz [Zhao
et al. 2001] y Pastry [Rowstron y Druschel 2001]. Estos sistemas abordan
Preocupaciones que no surgen para Bigtable, como el ancho de banda altamente variable,
participantes no confiables o reconfiguración frecuente; control descentralizado y
La tolerancia a fallas bizantinas no es un objetivo de Bigtable.
En términos del modelo de almacenamiento de datos distribuidos que se podría proporcionar
Para los desarrolladores de aplicaciones, creemos en el modelo de par clave­valor proporcionado por

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido para datos estructurados ∙ 4:23

Los árboles B distribuidos o las tablas hash distribuidas son demasiado limitantes. Los pares clave­
valor son un componente útil, pero no deberían ser el único que se proporcione a los desarrolladores.
El modelo que elegimos es más completo que los pares clave­valor simples y admite datos
semiestructurados dispersos. Sin embargo, sigue siendo lo suficientemente simple como para permitir
una representación de archivos planos muy eficiente, y es lo suficientemente transparente (mediante
grupos de localidades) como para permitir a nuestros usuarios ajustar comportamientos importantes
del sistema.
Varios proveedores de bases de datos han desarrollado bases de datos paralelas que pueden
almacenar grandes volúmenes de datos. La base de datos Real Application Cluster de Oracle
[[Link]] utiliza discos compartidos para almacenar datos (Bigtable utiliza GFS) y un gestor de
bloqueos distribuido (Bigtable utiliza Chubby). DB2 Parallel Edition de IBM [Baru et al. 1995] se basa
en una arquitectura de no uso compartido [Stonebraker 1986] similar a Bigtable. Cada servidor DB2
es responsable de un subconjunto de las filas de una tabla, que almacena en una base de datos
relacional local. Ambos productos ofrecen un modelo relacional completo con transacciones.

Los grupos de localidades de Bigtable logran beneficios de compresión y rendimiento de lectura


de disco similares a los observados en otros sistemas que organizan datos en disco utilizando
almacenamiento basado en columnas en lugar de basado en filas, incluido C­Store [Abadi et al.
2006; Stonebraker et al. 2005] y productos comerciales como Sybase IQ [French 1995; [Link]],
SenSage [[Link]], KDB+ [[Link]] y la capa de almacenamiento ColumnBM en MonetDB/X100
[Zukowski et al. 2005]. Otro sistema que realiza particionamiento vertical y horizontal de datos en
archivos planos y logra buenos índices de compresión de datos es la base de datos Daytona de AT&T
[Greer 1999]. Nuestros grupos de localidad no admiten optimizaciones a nivel de caché de CPU, como
las descritas por Ailamaki et al. [2001].

La forma en que Bigtable utiliza memtables y SSTables para almacenar actualizaciones en tablets
es análoga a cómo el Árbol de Fusión Estructurado en Registros [O'Neil et al. 1996] almacena
actualizaciones en los datos de índice. En ambos sistemas, los datos ordenados se almacenan en
memoria intermedia antes de escribirse en el disco, y las lecturas deben fusionar los datos de la
memoria y el disco.
C­Store y Bigtable comparten muchas características: ambos sistemas utilizan una arquitectura de
no compartición y cuentan con dos estructuras de datos diferentes: una para escrituras recientes y
otra para almacenar datos de larga duración, con un mecanismo para transferir datos de un formato a
otro. Los sistemas difieren significativamente en su API: C­Store se comporta como una base de datos
relacional, mientras que Bigtable ofrece una interfaz de lectura y escritura de bajo nivel y está diseñada
para soportar miles de operaciones de este tipo por segundo y servidor. C­Store también es un "SGBD
relacional optimizado para lectura", mientras que Bigtable ofrece un buen rendimiento tanto en
aplicaciones de lectura como de escritura intensivas.

El balanceador de carga de Bigtable debe resolver algunos de los mismos problemas de equilibrio
de carga y memoria que enfrentan las bases de datos sin recursos compartidos (p. ej., Copeland et al.
[1988] y Stonebraker et al. [1994]). Nuestro problema es algo más sencillo: (1) no consideramos la
posibilidad de múltiples copias de los mismos datos, posiblemente en formatos alternativos debido a
vistas o índices; (2) permitimos que el usuario nos indique qué datos pertenecen a la memoria y cuáles
deben permanecer en el disco, en lugar de intentar determinarlo dinámicamente; (3) no tenemos que
ejecutar ni optimizar consultas complejas.

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:24 ∙ F. Chang y otros.

11. CONCLUSIONES

Hemos descrito Bigtable, un sistema distribuido para almacenar datos estructurados.


en Google. Los clústeres de Bigtable se han utilizado en producción desde abril de 2005 y
Pasamos aproximadamente siete años­persona en diseño e implementación antes
En agosto de 2006, más de sesenta proyectos utilizaban Bigtable.
A nuestros usuarios les gusta el rendimiento y la alta disponibilidad que ofrece Bigtable
implementación, y que pueden escalar la capacidad de sus clústeres mediante
simplemente agregando más máquinas al sistema a medida que cambian sus demandas de recursos
con el tiempo.
Dada la inusual interfaz de Bigtable, una pregunta interesante es qué tan difícil ha sido para nuestros
usuarios adaptarse a su uso. Los nuevos usuarios a veces...
No están seguros de cuál es la mejor manera de utilizar la interfaz de Bigtable, especialmente si están
acostumbrados a usar bases de datos relacionales que admiten transacciones de propósito general. Sin
embargo, el hecho de que muchos productos de Google las utilicen con éxito...
Bigtable demuestra que nuestro diseño funciona bien en la práctica.
Estamos en el proceso de implementar varias funciones adicionales de Bigtable,
Como la compatibilidad con índices secundarios y la infraestructura para crear Bigtables replicadas entre
centros de datos con múltiples réplicas maestras. También tenemos
comenzó a implementar Bigtable como un servicio para grupos de productos, de modo que cada uno...
Los grupos no necesitan mantener sus propios clústeres. Como nuestros clústeres de servicio...
A mayor escala, tendremos que lidiar con más problemas de compartición de recursos dentro de Bigtable.
mismo [Banga et al. 1999; Bavier et al. 2004].
Finalmente, hemos descubierto que existen ventajas significativas al construir nuestro
Nuestra propia solución de almacenamiento en Google. Hemos obtenido una gran flexibilidad al diseñar
nuestro propio modelo de datos para Bigtable. Además, tenemos control sobre la implementación de
Bigtable y la demás infraestructura de Google.
Del cual depende Bigtable, significa que podemos eliminar cuellos de botella e ineficiencias a medida
que surjan.

EXPRESIONES DE GRATITUD

Agradecemos a los revisores anónimos de TOCS y OSDI, nuestro pastor de OSDI.


A Brad Calder y David Nagle por sus detalladas sugerencias para mejorar.
sobre los borradores de este documento.

El sistema Bigtable se ha beneficiado enormemente de los comentarios de nuestros numerosos


usuarios de Google. Además, agradecemos a las siguientes personas por sus contribuciones a Bigtable:
Shoshana Abrass, Dan Aguayo, Sameer Ajmani, Dan
Birken, Xin Chen, Zhifeng Chen, Bill Coughran, Jeff de Vries, Bjarni Einars­son, Mike Epstein, Shane
"
Gartshore, Fred´ eric Gobry, Healfdene Goguen, Robert Griesemer, Orla Hegarty, Jeremy Hylton, Josh
Hyman, Nick Johnson,
Alex Khesin, Marcel Kornacker, Joanna Kulik, Alberto Lerner, Shun­Tak
Leung, Sherry Listgarten, Mike Maloney, Wee­Teck Ng, Abhishek Parmar,
David Petrou, Eduardo Pinheiro, Kathy Polizzi, Richard Roberto, Deomid
Ryabkov, Shoumen Saha, Yasushi Saito, Steven Schirripa, Cristina Schmidt,
Lee Schumacher, Hao Shang, Jordan Sissel, Charles Spirakis, Alan Su, Chris
Taylor, Harendra Verma, Kate Ward, Peter Weinberger, Rudy Winnacker,
Frank Yellin, Yuanyuan Zhao y Arthur Zwiegincew.

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

Bigtable: un sistema de almacenamiento distribuido para datos estructurados ∙ 4:25

REFERENCIAS

ABADI, DJ, MADDEN, SR, Y FERREIRA, MC 2006. Integración de la compresión y la ejecución en sistemas de bases
de datos orientadas a columnas. Actas de la Conferencia Internacional ACM SIGMOD sobre Gestión de Datos. ACM,
Nueva York.
AILAMAKI, A., DEWITT, DJ, HILL, MD, y SKOUNAKIS, M. 2001. Relaciones de tejido para el rendimiento de caché. The
VLDB J. 169–180.
BANGA, G., DRUSCHEL, P., Y MOGUL, JC 1999. Contenedores de recursos: Una nueva herramienta para la gestión
de recursos en sistemas de servidores. En Actas del [Link] Simposio sobre Diseño e Implementación de Sistemas
Operativos. 45–58.
BARU, CK, FECTEAU, G., GOYAL, A., HSIAO, H., JHINGRAN, A., PADMANABHAN, S., COPELAND, GP, Y WILSON,
WG 1995. Edición paralela de DB2. IBM Syst. J. 34, 2, 292–322.
BAVIER, A., BOWMAN, M., CHUN, B., CULLER, D., KARLIN, S., PETERSON, L., ROSCOE, T., SPALINK, T. y WAWRZONIAK, M.
2004. Soporte de sistemas operativos para servicios de red a escala planetaria. En Actas del [Link] Simposio sobre Diseño e
Implementación de Sistemas en Red. 253–266.

BENTLEY, JL Y MCILROY, MD 1999. Compresión de datos mediante cadenas comunes largas. En Datos
Conferencia sobre compresión. 287–295.
BLOOM, BH 1970. Compensaciones espacio­temporales en la codificación hash con errores permitidos. Commun.
ACM 13, 7, 422–426.
BURROWS, M. 2006. El servicio de bloqueo Chubby para sistemas distribuidos débilmente acoplados. En Actas del 7.º
Simposio USENIX sobre Diseño e Implementación de Sistemas Operativos. 335–350.
CHANDRA, T., GRIESEMER, R. Y REDSTONE, J. 2007. Paxos en vivo: una ingeniería
perspectiva. En Actas del PODC.
Chang, F., Dean, J., Ghemawat, S., HSieh, WC, Wallach, DA, Burrows, M., Chandra, T., Fikes, A. y Gruber, RE 2006. Bigtable: Un
sistema de almacenamiento distribuido para datos estructurados. En Actas del 7.º Simposio USENIX sobre Diseño e
Implementación de Sistemas Operativos. 205–218.

COMER, D. 1979. Árbol B ubicuo. Computing Surveys 11, 2 (junio), 121–137.


COPELAND, GP, ALEXANDER, W., BOUGHTER, EE, Y KELLER, TW 1988. Ubicación de datos en Bubba. En Actas de
la Conferencia Internacional ACM SIGMOD sobre Gestión de Datos. ACM, Nueva York, 99–108.

DEAN, J. Y GHEMAWAT, S. 2004. MapReduce: Procesamiento de datos simplificado en grandes clústeres. En Actas
del 6.º Simposio USENIX sobre Diseño e Implementación de Sistemas Operativos.
137–150.

DEWITT, D., KATZ, R., OLKEN, F., SHAPIRO, L., STONEBRAKER, M., Y WOOD, D. 1984. Técnicas de implementación
para sistemas de bases de datos en memoria principal. En Actas de la Conferencia Internacional ACM SIGMOD sobre
Gestión de Datos. ACM, Nueva York, 1–8.
DEWITT, DJ Y GRAY, J. 1992. Sistemas de bases de datos paralelas: El futuro del alto rendimiento
sistemas de bases de datos. Commun. ACM 35, 6 (junio), 85–98.
FRANCÉS, CD 1995. Las arquitecturas de bases de datos universales no funcionan para DSS. En Actas de la Conferencia Internacional ACM SIGMOD
sobre Gestión de Datos. ACM, Nueva York, 449–450.

GAWLICK, D. Y KINKADE, D. 1985. Variedades de control de concurrencia en la ruta rápida IMS/VS. Datos.
Bull. Ing. 8, 2, 3–10.
GHEMAWAT, S., GOBIOFF, H., Y LEUNG, S.­T. 2003. El sistema de archivos de Google. En Actas del XIX Simposio de
la ACM sobre Principios de Sistemas Operativos. ACM, Nueva York, 29–43.
GRAY, J. 1978. Notas sobre sistemas operativos de bases de datos. En Sistemas Operativos: Un Curso Avanzado.
Apuntes de Informática, vol. 60. Springer­Verlag, ACM, Nueva York.
GREER, R. 1999. Daytona y el lenguaje Cymbal de cuarta generación. En Actas de la Conferencia Internacional ACM
SIGMOD sobre Gestión de Datos. ACM, Nueva York, 525–526.
HAGMANN, R. 1987. Reimplementación del sistema de archivos Cedar utilizando registro y confirmación grupal.
En Actas del 11º Simposio sobre Principios de Sistemas Operativos. 155–162.

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.
Machine Translated by Google

4:26 ∙ F. Chang y otros.

HARTMAN, JH Y OUSTERHOUT, JK 1993. El sistema de archivos de red con rayas de cebra. En Actas del XIV Simposio
sobre Principios de Sistemas Operativos. ACM, Nueva York, 29­43.
[Link]. [Link]/products/[Link]. Página del producto.
LAMPORT, L. 1998. El parlamento a tiempo parcial. ACM Trans. Comput. Syst. 16, 2, 133–169.
MACCORMICK, J., MURPHY, N., NAJORK, M., THEKKATH, CA, Y ZHOU, L. 2004. Boxwood: Abstracciones como base para
la infraestructura de almacenamiento. En Actas del 6.º Simposio USENIX sobre Diseño e Implementación de Sistemas
Operativos. 105–120.
MCCARTHY, J. 1960. Funciones recursivas de expresiones simbólicas y su cálculo por máquina. Commun. ACM 3, 4 (abr.),
184–195.
O'NEIL, P., CHENG, E., GAWLICK, D. y O'NEIL, E. 1996. El árbol de fusión con estructura logarítmica (árbol LSM). Acta Inf.
33, 4, 351–385.
[Link]. [Link]/technology/products/database/clustering/[Link]. Página del producto.
PIKE, R., DORWARD, S., GRIESEMER, R., Y QUINLAN, S. 2005. Interpretación de datos: Análisis paralelo con Sawzall.
Revista de Programación Científica 13, 4, 227–298.
RATNASAMY, S., FRANCIS, P., HANDLEY, M., KARP, R. y SHENKER, S. 2001. Una red escalable y direccionable por
contenido. En Actas de SIGCOMM. ACM, Nueva York, 161–172.
ROWSTRON, A. Y DRUSCHEL, P. 2001. Pastry: Ubicación y enrutamiento de objetos escalables y distribuidos
para sistemas peer­to­peer a gran escala. En Proceedings of Middleware 2001, pp. 329–350.
[Link]. [Link]/products­[Link]. Página del producto.
STOICA, I., MORRIS, R., KARGER, D., KAASHOEK, MF, Y BALAKRISHNAN, H. 2001. Chord: Un servicio de búsqueda
escalable entre pares para aplicaciones de Internet. En Actas de SIGCOMM.
ACM, Nueva York, 149–160.
STONEBRAKER, M. 1986. El argumento a favor de la nada compartida. Datab. Eng. Bull. 9, 1 (mar.), 4–9.
STONEBRAKER, M., ABADI, DJ, BATKIN, A., CHEN, X., CHERNIACK, M., FERREIRA, M., LAU, E., LIN, A., MADDEN, S.,
O'NEIL, E., O'NEIL, P., RASIN, A., TRAN, N., Y ZDONIK, S.
2005. C­Store: Un SGBD orientado a columnas. En Actas de la 10.ª Conferencia Internacional sobre Bases de Datos de
Gran Tamaño. ACM, Nueva York, 553–564.
STONEBRAKER, M., AOKI, PM, DEVINE, R., LITWIN, W., Y OLSON, MA 1994. Mariposa: Una nueva arquitectura para datos
distribuidos. En Actas de la 10.ª Conferencia Internacional sobre Ingeniería de Datos. IEEE Computer Society Press, Los
Alamitos, CA, 54–65.
[Link]. [Link]/products/databaseservers/sybaseiq. Página del producto.
ZHAO, BY, KUBIATOWICZ, J., Y JOSEPH, AD 2001. Tapestry: Una infraestructura para la localización y el enrutamiento de
áreas extensas con tolerancia a fallos. Informe Técnico UCB/CSD­01­1141, División de Ciencias de la Computación,
Universidad de California, Berkeley. Abr.
ZUKOWSKI, M., BONCZ, PA, NES, N., Y HEMAN, S. 2005. MonetDB/X100: un sistema de gestión de bases de datos (DBMS)
en la caché de la CPU. IEEE Data Eng. Bull. 28, 2, 17–22.

Recibido en diciembre de 2006; revisado en abril de 2008; aceptado en abril de 2008

ACM Transactions on Computer Systems, Vol. 26, No. 2, Artículo 4, Fecha de publicación: junio de 2008.

También podría gustarte