0% encontró este documento útil (0 votos)
11 vistas15 páginas

Diseño de Aplicaciones Intensivas en Datos

El documento aborda el diseño de aplicaciones intensivas en datos, destacando la importancia de la fiabilidad, escalabilidad y mantenibilidad en sistemas de datos. Se discuten diferentes modelos de datos, lenguajes de consulta y estructuras de almacenamiento, incluyendo NoSQL, SQL, y modelos de grafos, así como sus ventajas y desventajas. Además, se analizan técnicas de optimización y el funcionamiento interno de bases de datos, enfatizando la necesidad de elegir el motor de almacenamiento adecuado según las necesidades de la aplicación.

Cargado por

Violeta Saravia
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 TXT, PDF, TXT o lee en línea desde Scribd
0% encontró este documento útil (0 votos)
11 vistas15 páginas

Diseño de Aplicaciones Intensivas en Datos

El documento aborda el diseño de aplicaciones intensivas en datos, destacando la importancia de la fiabilidad, escalabilidad y mantenibilidad en sistemas de datos. Se discuten diferentes modelos de datos, lenguajes de consulta y estructuras de almacenamiento, incluyendo NoSQL, SQL, y modelos de grafos, así como sus ventajas y desventajas. Además, se analizan técnicas de optimización y el funcionamiento interno de bases de datos, enfatizando la necesidad de elegir el motor de almacenamiento adecuado según las necesidades de la aplicación.

Cargado por

Violeta Saravia
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 TXT, PDF, TXT o lee en línea desde Scribd

---

id: Designing Data-Intensive Applications


aliases: []
tags: []
---

Referencias: [Link]
## Capítulo 1

El software de porquería que vas a escribir en la mayoría de los shops es *data-


intensive*, no *compute-intensive*. Ese software se hace juntando bloques de
funcionalidad común que ya existen (caches, dbs, etc.) que llamamos *data systems*.
Las tres preocupaciones más grandes de estos sistemas son:

- Reliability: Funcionar correctamente, incluso ante adversidades (de software,


hardware o humanas).
- Scalability: Poder lidiar con el crecimiento (en volúmen de datos, de tráfico o
complejidad) del sistema.
- Maintainability: Poder mantener y adaptar la lógica del sistema de manera
productiva entre varios ingenieros.

### Faults != Failures

Fault es cuando un componente se desvía del spec. Failure es cuando el sistema o


componente deja de funcionar. Usamos mecanismos de fault-tolerance [2] para evitar
que las fallas se vuelvan fracasos. Una manera es aumentar la confianza en el
manejo de fallas es induciéndolas a propósito [1], ya que la mayoría pueden
prevenirse con testeos sencillos.

[1][Link]
[2][Link]
repid=rep1&type=pdf&doi=76535d0d4464f321bb572ecf1ec4a7855405276a

Tipos de failure:
- Hardware failures: Discos duros, power outages, AWS desactivando instancias.
- Software failures: Slowdowns, corrupted data causing errors, cascading failures.
- Human errors: Se solucionan con APIs y admin interfaces bien diseñadas, donde sea
difícil hacer algo mal, buen telemetry y sandboxes para desarrollo y testeo.

### Load

Para calcular la escalabilidad de un sistema, primero tenemos que describir su


load, y específicamente sus *load parameters* (ratio de reads/writes, hit rate de
un cache, req/sec, etc.).

E.g.: Twitter scaling, p.12

### Performance

Una vez que tenemos el load, poder describir la performance, preguntándonos que
pasa al variar un parámetro manteniendo los recursos del sistema, o cuántos
recursos hacen falta para mantener la performance al subir un load parameter.

Podemos medir la performance de un sistema en *response time*, usando percentiles.


p50 nos muestra el tiempo de respuesta usual para los usuarios; los p95+ (llamados
tail latencies) suelen ser financieramente importantes (e.g. en Amazon). Los
percentiles se usan en los *SLOs* y *SLAs*.
El response time muchas veces se ve afectados por los queue times, por lo que es
importante que los tests de response time no ocurran sólo secuencialmente.

Ver algoritmos para calcular percentiles:


1. Graham Cormode, Vladislav Shkapenyuk, Divesh Srivastava, and Bojian Xu:
“[Forward Decay: A Practical Time Decay Model for Streaming
Systems]([Link] at _25th
IEEE International Conference on Data Engineering_ (ICDE), March 2009.
2. Ted Dunning and Otmar Ertl: “[Computing Extremely Accurate Quantiles Using t-
Digests]([Link] _github.com_, March 2014.
3. Gil Tene: “[HdrHistogram]([Link] _hdrhistogram.org_.
4. [Link]
way-you-think/

### Maintainability

Hacer software fácil de mantener involucra tres cosas: operability, simplicity y


evolvability [3].

Los equipos de ops se encargan de monitorear la salud de sistemas, restaurarlos,


rastrear problemas de performance o failures, actualizar el software, anticipar
problemas, establecer buenas prácticas, documentar los sistemas de la organización,
etc. Para facilitar estas tareas, los sitemas pueden hacer lo siguiente:

- Buen monitoreo, para tener visibilidad del compartamiento de un sistema.


- Soporte para automatización.
- Evitar la dependencia de máquinas individuales.
- Proveer documentación.
- Proveer buenos defaults.

[3] [Link]

### Referencias

* [Link]
system-reliability

## Capítulo 2 - Data Models and Query Languages

### NoSQL

Necesidades de NoSQL:
- Mayor escalabilidad de la que permiten las bases relacionales.
- Queries especializadas.
- Modelos más dinámicos de los que permite la restricción de los esquemas
relacionales.
- Object-Relation mismatch: La mayoría del software hoy es OOP, lo que agregar la
necesidad de una capa de traducción entre la DB y los datos en la aplicación. Para
subsanar este problema se pueden usar ORMs, codificar la información en columnas de
JSON, o usar un document store como Mongo.

El modelo JSON tiene tres ventajas: 1) reducir el impedance mismatch entre la


aplicación y el storage, 2) no tener un schema y, para usos limitados, 3) la
representación JSON tiene mejor localidad de acceso que un esquema multi-tabla.

### SQL vs NoSQL

- Si los datos de la aplicación ya tienen estructura de documento (i.e. arbol de


one-to-many, donde el árbol completo se carga junto), una DB relacional podría
llevar a *shredding* innecesario.
- Si la aplicación tiene relaciones many-to-many, los datos tiene que estar
denormalizados en el Document Model, o se tienen que emular joins con varias
requests, lo que agrega complejidad a la aplicación.
- Los Document DBs se suelen llamar schemaless, pero sería mejor llamarlas *schema-
on-read*, i.e., el esquema no es chequeado por la DB y sólo se interpreta al leer
los datos. Lo opuesto sería *schema-on-write*, donde el esquema es explícito y la
DB chequea que los datos estén conformes. Cada tiene ventajas y desventajas para
ciertos usos (startups vs. empresas, data heterogenea vs homogenea)
- *Localidad de datos*: Joins vs single document retrieval. Las DBs relacionales
también tienen maneras de agrupar datos relacionados: Spanner [4] permite anidar
filas de una tabla en una tabla padre. Oracle también, usando *multi-table index
cluster tables*, además del concepto *column-family* usando en las Bigtables [5]
(Cassandra, HBase, etc.).

[4] [Link]
[5] [Link]
structured-data/

### MapReduce

Aunque MapReduce ya no se usa, algunas Document Stores como MongoDB tienen


extensiones de MapReduce para hacer queries limitadas:

```javascript
[Link](
function map() {
var year = [Link]();
var month = [Link]() + 1;
emit(year + "-" + month, [Link]);
},
function reduce(key, values) {
return [Link](values);
},
{
query: { family: "Sharks" },
out: "monthlySharkReport"
}
);
```

MongoDB también soporta un lenguaje declarativo llamado *aggregation pipeline*:

```javascript
[Link]([
{ $match: { family: "Sharks" } },
{ $group: {
_id: {
year: { $year: "$observationTimestamp" },
month: { $month: "$observationTimestamp" }
},
totalAnimals: { $sum: "$numAnimals" }
} }
]);
```

### Conclusión

Argumentos para nunca usar Document Stores:


- Las relaciones siempre van a ser relevantes eventualmente.
- Las pocas ventajas de los Document Stores ya fueron absorbidas por las DBs
relacionales (Postgres JSON column, etc.).
- Es más fácil ir de relacional a document que lo opuesto.
- MongoDB terminó volviendose SQL.

### Modelos de Grafos

Vértice: ID, edges entrantes, salientes, y propiedades (k:v pairs)


Edge: ID, tail y head vertex, label y propiedades (k:v pairs)

### Cypher

```Cypher
CREATE
(NAmerica:Location {name:'North America', type:'continent'}),
(USA:Location {name:'United States', type:'country' }),
(Idaho:Location {name:'Idaho', type:'state' }),
(Lucy:Person {name:'Lucy' }),
(Idaho) -[:WITHIN]-> (USA) -[:WITHIN]-> (NAmerica),
(Lucy) -[:BORN_IN]-> (Idaho)
```

```Cypher
// Find people who emigrated from the US to Europe
MATCH
(person) -[:BORN_IN]-> () -[:WITHIN*0..]-> (us:Location {name:'United
States'}),
(person) -[:LIVES_IN]-> () -[:WITHIN*0..]-> (eu:Location {name:'Europe'})
RETURN [Link]
```

### Triple-Stores

Similar a los property graphs, con algunas herramientas distintas. Los datos se
representan con triples *(subject, predicate, object)*, e.g.:

```
@prefix : <urn:example:>.
_:lucy a :Person.
_:lucy :name "Lucy".
_:lucy :bornIn _:idaho.
_:idaho a :Location.
_:idaho :name "Idaho".
_:idaho :type "state".
_:idaho :within _:usa.
_:usa a :Location.
_:usa :name "United States".
_:usa :type "country".
_:usa :within _:namerica.
_:namerica a :Location.
_:namerica :name "North America".
_:namerica :type "continent".

// Equivalente:
@prefix : <urn:example:>.
_:lucy a :Person; :name "Lucy"; :bornIn _:idaho.
_:idaho a :Location; :name "Idaho"; :type "state"; :within _:usa.
_:usa a :Location; :name "United States"; :type "country"; :within _:namerica.
_:namerica a :Location; :name "North America"; :type "continent".
```

Los vértices se escriben *_:someName*. Los predicados son edges si el objeto es un


nodo, o keys si el objeto es un primitive.

### SPARQL

```SPARQL
PREFIX : <urn:example:>

# Equivalente a la query de Cypher


SELECT ?personName WHERE {
?person :name ?personName.
?person :bornIn / :within* / :name "United States".
?person :livesIn / :within* / :name "Europe".
}
```

Lo interesante del formato RDF es que no distingue entre propiedades y edges, ambos
son predicatos. La misma sintáxis se puede usar para buscar propiedades o tipos de
edges.

## Capítulo 3 - Storage and Retrieval

Para que es necesario entender el funcionamiento interno de una DB? Porque hay
muchos, y hace falta encontrar el apropiado para una aplicación, además de para
poder tunearlo para nuestro workload. Por ejemplo, hay mucha diferencia entre los
storage engines optimizados para transacciones y los optimizados para analytics.

### Indexes

Un índice es una estructura de datos adicional, derivada de los datos insertados en


la DB; pueden agregarse y removerse sin afectar los datos. Mantener indices
ralentiza los writes, pero optimiza mucho los reads, por eso no suelen estar
incluídos por default.

### Hash Indexes

Si tenemos un log store (i.e. una secuencia de records append-only) podemos indexar
la tabla con k-v pairs que almacenan el index y un pointer a su ubicación en
memoria.

Ventajas de los logs append-only:


- Permiten llevar a cabo writes con writes secuenciales.
- Hacen mucho más sencilla la concurrencia y recuperación de crashes.

Desventajas del hash index:


- La tabla tiene que entrar en RAM para ser eficiente.
- Las queries por rango no son eficientes: cada item requiere un hash individual, y
los datos están en lugares random del store.

### SSTables y LSM-Trees

Sorted String Table: Log table en la que las secuencias de k-v están ordenadas por
key. Ventaja: No es necesario mantener el index completo en memoria. Podemos
mantener partes del index y escanear entre keys para encontrar un valor específico.
Para mentener los writes entrantes ordenados por key usamos una estructura de datos
auto-organizada como los R-B trees. Finalmente, el storage engine funciona así:
- Los writes entrantes se agregan a una estructura balanceada en memoria, a veces
llamada *memtable*.
- Al superar cierto tamaño, los datos se escriben a un archivo SSTable, que se
vuelve el segmento más reciente de la DB. Los writes entrantes pueden continuar en
una memtable nueva sin problemas.
- Al entrar un read, se busca la key en la memtable, luego en el segmento más
reciente, etc.
- Un proceso periódico mergea y compacta los archivos de segmento en el disco,
descartando valores sobreescritos o eliminados.

Problema: La memtable se pierde en caso de un crash. Para evitarlo, mantenemos un


log separado en el disco donde se apende inmediatamente cada write. Este log no
está organizado ni es eficiente, pero no importa, ya que sólo existe para recuperar
la memtable en caso de un crash, y al guardarse a una SSTable, se descarta.

Este algoritmo se usa en LevelDB, RocksDB, además de Cassandra y HBase.


Originalmente se conocía como *Log-Structured Merge-Tree*.

### Optimizaciones

Para evitar muchos reads innecesario cuando la key no existe en la tabla, se pueden
usar bloom filters para chequear si está presente antes de intentar hashearla.

### B-Trees

Los B-Trees dividen a la base de datos en *blocks* o *pages* de tamaño fijo,


típicamente 4KB, y leer o escriben una página a la vez. Cada página se identifica
con una dirección; al mantener las keys organizadas, cada read se vuelve una
búsqueda binaria al descender el árbol.

El número de referencias en cada página no-hoja es el *branching factor*. Depende


del espacio requerido para almacenar las referencias, pero suele ser de varios
cientos.

Durante un write, si el record nuevo no entra en la página correspondiente, esta se


divide en dos y se crea una nueva página que refiere a ambas. Esto significa que
los writes a veces requieren de vario writes en el disco; una operación peligrosa
que puede dejar un indice corrompido si la DB crashea entre un write y otro. Para
evitar esto, las implementaciones de B-Trees suelen incluir un *write-ahead log*
(WAL) o *redo log*, un archivo append-only donde se escriben todas las
modificaciones de un B-Tree antes de aplicarse. En caso de un crash, puede usarse
este log para restaurar el B-Tree.

### Optimizaciones

- En vez de mantener un WAL, se puede hacer copy-on-write, escribiendo cada página


modificada a una ubicación distinta, y creando una nueva versión de las páginas
parent apuntando a esta ubicación nueva.
- Las páginas pueden estar en cualquier lugar del disco, lo que puede hacer
ineficientes las range-searches. Muchos B-Trees intentar mantener las hojas
ordenadas secuencialmente en el disco, lo cual es dificil de mantener al crecer la
DB.
- Los *fractal trees* toman ideas de los logs para reducir accesos al disco. [6]

[6] Bradley C. Kuszmaul: “A Comparison of Fractal Trees to Log-Structured Merge


(LSM) Trees,” [Link], April 22, 2014.

### Ventajas de los LSM-Trees


*Write amplification*: El efecto de un write a la DB resultando en multiples writes
en el disco. En los B-Trees, ocurre por los WALs, además de que cada página se
sobreescribe entera al tener que hacer un cambio; algunos storage engines incluso
sobreescriben la misma página dos veces para evitar con updates parciales en cada
de una falla de energía. En los LSMs, ocurre por el compactamiento y mergeo de
SSTables.

- En las aplicaciones write-intensivas, los LSM suelen tener ventajas por 1) poder
sostener un throughput de writes mayores, dependiendo de la configuración del
engine, ya que su write amplification es menor, y 2) poder escribir secuencialmente
archivos SSTable en lugar de tener que sobreescribir páginas de un B-Tree,
particularmente en HDDs.
- Los datos se representan de manera más compacta, y pueden comprimirse mejor, que
en los B-Trees, lo que permite más reads y requests por I/O, incluso en SSDs.

### Desventajas de los LSM-Trees

- El write I/O debe compartirse entre las requests y el proceso de compactamiento.


Aunque los engines intenten optimizar esto, hay percentiles altos en los que el
response time de las queries a un LSM puede ser muy alto, ya que la request debe
esperar a que el disco termine una compactación cara. Los B-Trees tienen
performance más predecible.
- Cuando más grande una DB, más I/O del disco es necesario para la compactación. Si
el write throughput es alto y la compactación está mal configurada (o no es
monitoreada) [7][8], los writes pueden exceder la velocidad de compactación. Esto
hace que los segmentos sin mergear crezcan sin límite superior hasta que el espacio
en disco se acabe, además de ralentizar mucho los reads.

[7] Benjamin Coverston, Jonathan Ellis, et al.: “CASSANDRA-1608: Redesigned


Compaction, [Link], July 2011.
[8] Igor Canadi, Siying Dong, and Mark Callaghan: “RocksDB Tuning Guide,”
[Link], 2016.

- En los B-Trees, cada key existe exactamente en un lugar del index, lo que lo hace
más útil para ofrecer semánticas transaccionales más fuertes (ver cap. 7).

### Índices Secundarios

Los índices secundarios son cruciales para, por ejemplo, llevar a cabo join de
manera eficiente, indexando foreign keys. La key en el index es lo que usan las
queries para buscar, pero el valor puede estar en dos lados: en la fila en
cuestión, o referenciada en otro lado (*heap file*).

Los heap files evitan duplicar data cuando hay vario índices secundarios, además de
optimizar los updates: si el nuevo valor es igual o menor al anterior, puede
cambiarse en el lugar. De lo contrario, hay que actualizar todos los índices, o
dejar un forwarding pointer en la ubicación anterior y encontrar un espacio nuevo
en el disco.

Si el salto del índice al heap file afecta mucha la performance de los reads, la
fila indexada puede almacenarse directamente en un índice, un *clustered index*. La
primary key de una tabla en InnoDB siempre es un clustered index, y los índices
secundario refieren a la primary key en lugar de a un heap file. En MSSQL se puede
especificar un clustered index por tabla [9].

[9] Books Online for SQL Server 2012. Microsoft, 2012.

El intermedio entre un clustered index y heap files es un *covering index* o *index


with included columns*, donde algunas columnas de una tabla se almacenan con el
índice y otras no, optimizando algunas queries [10].

[10] Joe Webb: “Using Covering Indexes to Improve Query Performance,” simple-
[Link], 29 September 2008.

### Índices Multi-Columna

El más común se llama *índice concatenado*, que combina varios campos en una key
anexando una columna a otra en un orden determinado.

Los índices multidimensionales son importantes, por ejemplo, para los datos
geoespaciales:

```SQL
SELECT
*
FROM
restaurants
WHERE
latitude > 51.4946 AND
latitude < 51.5079 AND
longitude > -0.1162 AND
longitude < -0.1004;
```

Un índice B-Tree o LSM común no podría recorrer esta query de manera eficiente, ya
que tendría que indexar la latitud *o* la longitud, pero no ambas. La solución más
común es usar una estructura especializada como los R-Trees, como hace PostGIS en
PostgreSQL. [11]

[11] The PostGIS Development Group: “PostGIS 2.1.2dev Manual,” [Link], 2014.

Lo índices multi-dimensionales no son sólo para ubicaciones geográficas. Por


ejemplo, en una web de ecommerce se podría usar un índice 3D en (red, green, blue)
para buscar productos en un rango de colores. Una DB climática podría indexar
(date, temp) en 2D para buscar observaciones de un año X con un rango de
temperatura Y.

### Fuzzy Indexes

Para buscar keys *similares* a otras, e.g. al hacer text search, se pueden usar
índices que agrupen keys a cierta distancia de edición. En Lucene, por ejemplo, el
index en-memoria es similar a un trie; puede convertirse en un Levenshtein
automaton, que soporta la búsqueda de palabras dentro de una distancia de edición.
[12]

[12] Klaus U. Schulz and Stoyan Mihov: “Fast String Correction with Levenshtein
Automata,” International Journal on Document Analysis and Recognition, volume 5,
number 1, pages 67–85, November 2002. doi:10.1007/s10032-002-0082-8

Otras técnicas de fuzzy search incluyen machine learning y clasificación de


documentos. [13]

[13] Christopher D. Manning, Prabhakar Raghavan, and Hinrich Schütze: Introduc‐


tion to Information Retrieval. Cambridge University Press, 2008. ISBN: 978-0-521-
86571-5, available online at [Link]/IR-book

### In-Memory Stores


Las DBs In-Memory ofrecen ventajas de performance, guardando snapshots al disco
sólo por durabilidad, pero sirviendo los reads enteramente desde RAM. Parte de la
ventaja que tienen es que no necesitan serializar los datos de una manera que pueda
escribirse al disco. Esto, además, les permite ofreces modelos de datos difíciles
de implementar con indexes en disco, como los priority queues y sets que ofrece
Redis.

### OLTP vs OLAP

| | OLTP | OLAP |
| --------- | ---------------------- | ------------------------- |
| Reads | Pocos records, por key | Muchos records, agregados |
| Writes | Random, low-latency | Bulk (ETL) o event stream |
| Usado por | Webapps | Internal analysts |
| Datos = | Presente | Historial |
| Tamaño | Gb/Tb | Tb/Pb |

### Data Warehousing

Para evitar que las consultas de analytics arruinen la latencia de un frontend, se


usan DBs separadas llamadas *data warehouses*, que contienen una copia read-only de
los OLTPs de una empresa. Los datos pueden transformarse en un schema más amistoso
para el análisis, limpiarse y luego cargarse al data warehouse (ETL).

La ventaja de los data warehouses es que se puede optimizar el acceso. Los


algoritmos de indexado que vimos hasta ahora no son tan buenos para queries
analíticas.

Algunas DBs como MS SQL y SAP HANA soportan tanto OLTP como data warehousing, pero
en general se están volviendo dos productos separados.

### Schemas para Analytics

El modelo de datos para data warehouses es bastante formuláico, conocido como *star
schema* o *dimensional modeling*. En el centro está la *fact table*, con una fila
para cada evento que haya ocurrido. Algunos columns son foreign keys a otras tabla,
llamadas *dimension tables*. Como la fact table es un evento, las dimensiones son
el *quien*, *donde*, *por qué*, etc. También se puede normalizar más en un
*snowflake schema* (lol) pero generalmente se prefieren los stars porque son más
sencillos de analizar.

### Columnar Storage

En las DBs OLTP, sean relacionales o document stores como los que vimos hasta
ahora, los datos se organizan por filas: todos los valores de una fila de la tabla
están uno al lado del otro. Para procesar queries de analytics generalmente
necesitamos buscar rangos grandes de ciertas columnas, sin devolver la fila entera
(que en un data warehouse, estando no normalizado, puede ser muy grande).

Almacenar los datos por columna también los hace más fáciles de comprimir, ya que
los patrones repetitivos de los datos están ordenados en secuencia. Los datos
también son más fáciles de vectorizar y utilizan mejor el cache de la CPU.

### Sort Order en Columnar Storage

Dado que las columnas deben organizarse de manera igual (para que sepamos a qué
fila pertenece cada elemento), sólo podemos elegir una secuencia de columnas para
organizar los datos, basándonos en queries comúnes.
Si queremos distintos sort orders, podríamos duplicar los datos, ordenándolos de
maneras distintas, aprovechando que los datos tienen que ser replicados por
seguridad de todos modos.

### Aggregation

Dado que las data warehouses se usan para analytics, y muchas cálculos de analytics
se repiten (COUNT, SUM, AVG, etc.) podemos cachearlos en un *materialized view*. A
diferencia del view de una base relacional, que es sólo una tabla "virtual" que
ejecuta una query al leerse, una materialized view es una copia de los resultados
de la query escritos en el disco.

Un caso especial común de materialized view es lo que se llama un *data cube* o


*OLAP cube*, una grilla de agregados agrupados en dimensiones distintas. Para *N*
dimensiones (i.e. tablas) podemos imaginar un cubo de *N* dimensiones con los
agregados de una combinación particular de valores en cada dimensión. La desventaja
es que los data cube tienen mucha menos flexibilidad, por lo que al final y al cabo
siempre hace falta almacenar los datos en tablas, y usar cubos sólo para optimizar
ciertas queries.

## Capítulo 4 - Encoding and Evolution

Los cambio en el formato de datos o schema suelen requerir de un cambio en el


código. Pero el código no puede cambiar instantaneamente: en los servidores grandes
se suelen hacer rolling upgrades, y en el client-side dependés de que el usuario
actualice la versión. Para que esto funcione las versiones nuevas y viejas deben
tener *backward* y *forwards compatibility*.

### Formatos de Encoding

Las datos se trabajan en al menos dos representaciones distintas:


- En memoria (structs, lists, hash tables, etc.)
- Encoded (JSON, XML, etc.)

### Text Formats

Los formatos de texto (JSON, CSV, etc.) tienen varios problemas:


- El encoding de numbers es ambiguo. CSV y XML no distinguen strings con números de
números, y JSON no distingue int, float32 y float64. Al lidiar con números grandes,
como las IDs de tweets, la API de twitter devuelve la ID dos veces, una como número
y otra como string para JavaScript, que no tienen ints y podría interpretar números
muy grandes mal.
- JSON y XML soportan unicode strings, pero no binary strings. Para subsanar esta
limitación se puede codificar la data binaria como texto con Base64 (esto aumenta
el tamaño de los datos un 33%).
- XML y JSON tienen schema, pero CSV no, y no todos los parsers implementan
detalles del formato correctamente (e.g. comas).

### Binary Formats

Cuando la comunicación es interna y los datasets son muy grandes (terabytes), puede
ser util enviar los datos en un formato binario (como RDB!). *Thrift* y *protobuf*
definiendo primero un schema, procesado por una herramienta de codegen.

#### Protobuf

```protobuf
message Person {
required string user_name = 1; // 1 == field tag
optional int64 favorite_number = 2;
repeated string interests = 3;
}
```

Notice that in the case of datatype changes, two things can happen:
- A change in size may cause a loss of precision if you, e.g., go from 32 to 64
bits.
- In the case of protobuf, there isn't an array type, but instead a `repeated` tag,
so that a schema change from `optional` or `required` would only cause the last
value to be read in one schema, and all of them in the other, ensuring
compatibility.

#### Avro

```avro
record Person {
string userName;
union { null, long } favoriteNumber = null;
array<string> interests;
}
```

Este formato es mucho más compacto que los otros, porque el bytecode no almacena
tags ni ninguna otra información: sólo los bytes, un dato tras otro. Esto significa
que un binario avro sólo se puede decodificar con el schema exacto. Cómo soportar
schema evolution, entonces?

En Avro, los datos se codificar en un schema (el *writer's schema*) y se


decodifican en otro que tenga la aplicación que los lee (el *reader's schema*).
Estos dos no necesitan ser iguales; al decodificar, la librería Avro resuelve las
diferencias mirando ambos.

Para mantener compatibilidad, sólo se pueden agregar o quitar campos con valores
default. Los nombres se pueden cambiar manteniendo aliases, y los tipos sólo si se
pueden convertir.

Ok, pero, dónde está el writer's schema? Depende del contexto en que se use Avro:
- Archivo grande con muchos registros: Se puede incluir al comienzo del archivo,
sin perder muchos space savings.
- DB con registros individuales: Se puede incluir un número de version en cada
record, y mantener una lista de versiones del schema (ver Espresso).
- Records enviados por red: Al comunicarse por una red bidireccional, los procesos
pueden acordar un schema al establecer la conección. El protocolo Avro RPC funciona
así.

La ventaja de Avro es que, al no tener tag numbers en el schema, es mucho más fácil
trabajar con *dynamically generated* schemas. Si el schema cambia, sólo hay que
generar otro. En protobuf o Thrift, en cambio, tendrías que reasignar field tags a
mano, cuidándote de no reusar tags usados anteriormente. Además, en lenguajes
dinámicos, Avro no necesita necesariamente usar un API generada con codegen. Puede
abrir el filo y, si es *self-describing*, leerlo por si mismo. Apache Pig funciona
así.

Además de los mencionados, hay encodings binario propietarios, como los usados para
las queries de las DBs relaciones (a través de las APIs ODBC o JDBC, e.g.).

### Dataflows Reales


#### DBs

Si una versión vieja de la aplicación intenta deserializar una tabla con un campo
nuevo y volverlo a enviar modificado, ese campo debería mantenerse intacto. Esto es
fácil de implementar pero hay que tener cuidado de no perder el campo en el
proceso.

Para archivamiento, se guardan los datos junto con su schema. Dado que son
inmutables, es una buena idea almacenarlos en Avro containers o en un formato
columnar como Parquet.
#### Web Services

Ver [Link]

Tipos de servicios:
- REST: Formatos de datos sencillos, recursos identificados por URL, uso de
propiedades de HTTP para auth, content type, cache control, etc. usando un formato
de definición como OpenAPI.
- SOAP: XML-based, independiente de HTTP, viene con enormes estándares relacionados
(los *web service frameworks*, WS-\*). Formatos y datos no legibles, difíciles de
integrar por organizaciones externas.
- Viejos: Muchos protocolos RPC viejos, incluído SOAP, intentan parecerse lo más
posibles al llamado de una función local (de ahí el nombre RPC). Esto es muy malo,
obviamente. Las requests a través de la red son fundamentalmente distintas de una
función local (más failure states, encoding/decoding de datos que podrían ser de
otro lenguaje de programación, etc.). Por eso REST es mejor.
- Nuevos: La nueva generación de RPCs (de protobuf (gRPC), Avro, Thrift, etc.)
explicitan más estos estados posibles, e.g. usando *futures* para encapsular
acciones asincrónicas. gRPC tiene *streams*, donde una llamada consiste de un
iterador de requests y responses a lo largo del tiempo. Estos frameworks también
soportan *service discovery*.

#### Async Message Passing

#### Message Brokers

Más comúnes: RabbitMQ, ActiveMQ, NATS, Apache Kafka.

Ventajas:
- Buffer para recipientes no disponibles o sobrecargados.
- Redeliver automático.
- El receiver no necesita conocer la dirección del sender.
- One-to-Many messages
- Desacople lógico del sender y el reciever (ponele).

Deventajas:
- Más lento.
- One-way dataflow
- Si un consumer republica un mensaje a otro topic, tiene que tener cuidado de
preservar campos desconocidos.

#### Distributed Actor Frameworks

Los *actos* son patron de diseño para concurrencia. Cada actor es un cliente o
entidad que se comunica con otros actos enviando mensajes asincrónicos. El delivery
de mensajes no está garantizado. Como cada actor procesa un mensaje a la vez, no
hace falta preocuparse por threads; cada actor se puede coordinar
independientemente por el framework.
## Parte II - Distributed Data

Por qué distribuir una DB en varias máquinas?


- Scalability: Si el volume de datos, reads o writes, excede lo que puede manejar
una sola máquina.
- Fault tolerance/availability.
- Latency: Datacenters locales para distintas partes del mundo.

### Scaling

El scaling más sencillo, el vertical, involucra simplemente agregar memoria, disco


o CPUs para aumentar el load posible. El problema es que en la arquitecturas
*shared-memory* el precio no escala.

En las arquitecturas *shared-nothing* (también llamado *horizontal scaling*) cada


PC es un nodo, usando sus recursos independientemente, de dos maneras distintas:

- Replication: Mantener copias redundantes de los datos en distintos nodos.


- Partitioning: Dividir la DB en particiones asignadas a nodos distintos. También
llamado *sharding* (ver Citus en Postgres).

## Capítulo 5 - Replication

Algoritmos para replicar cambio entre nodos: *Single-leader, multi-leader y


leaderless*. También hay otros tradeoffs como sync/async y manejo de fallas.

### Leader-based replicas

1) Una de las réplicas se designa como líder. Los clientes enviar requests de
writes al líder.
2) Las otras réplicas son *followers*. Cuando el líder escribe cambios, los reenvía
a sus followers en un *relication log* o *change stream*. Los followers toman el
log y ejecutan los writes en el mismo orden. Los conceptos de hot, warm y cold
standbies también se usan pero con significados varios: En Postgres, *hot standby*
es una réplica que acepta reads de clientes y *warm standby* procesa cambio del
líder, no de los clientes.
3) Los reads de los clientes se pueden procesar en el líder o en los followers, los
writes sólo en el líder.

Ejemplos: Postgres 9+, SQL Server's AlwaysOn Availability Groups. También se usan
en DBs no relacionales y en brokers como Kafka o RabbitMQ.

### Sync vs Async

En una réplica *sync*, el líder esperan a que su réplica reporte OK en la


transacción para confirmar la operación al usuario. Así, ese follower siempre tiene
garantizadas copias actualizadas de los datos. Si el líder falla, siempre podemos
estar seguros de que el follower tiene los datos. La desventaja es que la operación
no se procesa, o es lenta, si el follower falla o tarda. En la práctica se usan
tanto sync como async, lo que se llama una configuración *semi-synchronous*.

Si bien la replicación *async* no garantiza durabilidad de los datos (si el lider


falla y no es recuperable), se suele usar bastante, ya que suele haber muchos
followers o estar geográficamente distribuidos.

### Configurar followers

Cómo configurar followers sin downtime:


1) Tomar un snapshot de la DB, preferentemente sin lockear la DB entera. La mayoría
de las DBs tiene esta feature, ya que se usa para backups.
2) Copiar el snapshot al follower.
3) El follower se conecta al líder y solicita los cambios desde que se tomo la
snapshot. Esto requiere que la snapshot esté asociada a una posición exacta en el
replication log del líder (*log sequence number/LSN* en PostgreSQL y SQL Server,
*binlog coords* en MySQL).
4) El follower se puso al día.

### Fallas en Nodos

Falla de líder:
1) Se determina que el líder fallo con un timeout.
2) Se elige un nuevo líder, ya sea un *controller node* previamente asignado, o el
follower más actualizado, etc.
3) Se reconfiguran los followers para que apunten al nuevo líder, y cuando el líder
vuelva, se transforma en un follower.

Problemas:
- Con replicas *async*, el lider viejo podría quedar con writes que los otros no
tienen. La solución más fácil es descartarlos. Esto es particularmente peligroso si
hay otras DBs externas que necesitan los mismos contenidos [14].
- Dos nodos podrían creer que son líderes. Si ambos están aceptando writes, podría
haber conflictos sin resolución posible.
- Cuál es el timeout correcto? Uno largo hace más lenta la recuperación, pero uno
corto puede causar failovers innecesarios que empeoren más aún problemas en la red.

En la práctica, algunos equipos prefieren llevar a cabo failovers manualmente.

### Replication Logs

- Statement-based: Enviar las queries a los followers. Sólo funciona con queries
determinísticas.
- WAL shipping: Enviar el write-ahead log del líder a los followers. El problema es
que los WALs almacenan información de bajo nivel (i.e. los bytes cambiados en tal o
cual bloque) por lo que un mismatch de versiones de la DB entre líder y follower
podría volverlos inútiles; si los WALs se usan como log, una update de la DB
requiere bajar todos los nodos.
- Logical (row-based) replication: Decoplar los logs de replicación y storage,
solucionando los problemas de arriba.
- Trigger-based: Si necesitás una replicación más customizada (replicas sólo un
subset de datos, o de un tipo de DB a otra, etc.) se puede pasar la replicación a
la aplicación, o usando *triggers* o *stored procedures* de la DB para ejecutar un
logeo en una tabla separada durante una transacción. Luego un proceso externo puede
aplicar las transformaciones necesarias y replicas los datos en otro sistema (e.g.
*Bucardo* en Postgres o *Databus* en Oracle).

### Soluciones al Replication Lag

- *Read-after-write Consistency*: Para evitar que un usuario no pueda ver los


cambios que acaba de hacer, hacemos que los reads sean desde el líder cuando sea
necesario: Cuándo el usuario modifica su propio perfil, o durante una ventana de
tiempo luego del ultimo update enviado, o guardando un timestamp del último write y
leyendo sólo de nodos que no estén atrasados al timestamp.
- *Monotonic Reads*: Evitar que varios reads de un usuario sean a nodos con lags
distintos, por ejemplo, haciendo que cada usuario lea siempre de la misma réplica,
por ejemplo, asignando un usuario a una réplica con un hash.
- *Consistent Prefix Reads*: Garantizar que, si una secuencia de writes ocurre en
un cierto orden, cualquiera que los lea los vea aparecer en el mismo orden. Por
ejemplo, para que los mensajes de un chat no aparezcan desordenados.
Particularmente problemático en DBs particionadas; si los writes se aplican en el
mismo orden, los reads siempre van a tener un prefijo consistente.

También podría gustarte