Parciales viejos: [Link] ([Link]
edu/courses/cs/186/)
03 SQL II
Cross Join
FROM tabla1, tabla2 -- CROSS JOIN entre tabla 1 y 2
Set Comparators
[...] EXISTS columna1 -- columna1 is not empty
-- (equivalente a IS NOT NULL para columnas?)
-- También están los comparadores de conjuntos
-- (NOT) IN, (NOT) EXISTS, op ANY, op ALL
-- Number of sailors who've reserved boat #102
SELECT [Link]
FROM Sailors S
WHERE EXISTS
(SELECT *
FROM Reserves R
WHERE [Link] = 102 AND [Link] = [Link])
-- El uso de [Link] hace que la subquery sea *correlacionada*,
-- por lo que la subquery se recomputa para cada tupla de S.
-- Como si la subquery fuese una función a la que se le pasa cada fila S
"Division"
SELECT [Link]
FROM Sailors S
WHERE NOT EXISTS
(SELECT [Link]
FROM Boats B
WHERE NOT EXISTS
(SELECT [Link]
FROM Reserves B
WHERE [Link] = [Link]
AND [Link] = [Link]))
ArgMax
-- Single Max (not deterministic,
-- doesn't account for multiple maxes)
SELECT S.*
FROM Sailors S
ORDER BY [Link] DESC
LIMIT 1;
-- ARGMAX
SELECT *
FROM Sailors S
WHERE [Link] >= ALL
(SELECT [Link]
FROM Sailors S2)
Views
Lazy-evaluated saved queries. También sirven para permisos: Se puede dar acceso a un view y no a sus tablas subyacentes.
CREATE VIEW <view_name>
AS <select_statement>
Common Table Expression (CTE)
Definimos una subquery antes del statement para reutilizarla en otras queries.
WITH <identifier> AS <select_statement>, *
<select_statement>
NULL
Todos los tipos de datos pueden ser NULL
Los outer joins los producen naturalmente cuando no hay match entre tablas
Checks con NULLs evalúan a falso
Agregados con NULL, en líneas generales, ignoran el NULL.
04 Disks, Files and Buffers
DBMS Parsing
Trayecto SQL Client -> Database:
1. Query Parsing and optimization
2. Relational Operators: algoritmos individuales que se arman para formar un dataflow
3. Files and Index Management: Organiza tablas y registros como grupos de páginas en un archivo lógico
4. Buffer Management: Provee la ilusión de estar trabajando en RAM. Las capas de arriba asumen que están trabajando
en memoria, y el buffer manager se encarga de mapear bloques del HDD a RAM.
5. Disk Space Management: Traduce page requests del buffer a bytes físicos en el almacenamiento.
Cada capa abstrae las capas inferiores.
Dos módulos/problemas cortan todo el trayecto: Concurrency control y recovery.
Discos
Los sistemas de bases de datos no nos permiten interactuar con datos directamente, al final de bytes y pointers. En su lugar,
hay una API que lee y escribe "páginas" de datos del disco a ram (READ) y de RAM al disco (WRITE). Esto minimiza los
errores comunes de pointers que ocurren en C.
06 Índices
Un índice es una estructura de datos que acelera la búsqueda y modificación de entradas por clave: El lookup puede ser por
igualdad, por rango o región, etc. y la búsqueda puede aceptar cualquier subconjunto de columnas en la relación.
IBM Indexed Sequential Access Method (ISAM)
1. Para un heap file ordenado y almacenado de manera contigua (i.e. no hacen falta pointers a las páginas) queremos
construirle un índice. No usamos binary search sobre la base porque, para una base muy grande, necesitamos
muchos IOs, \(log_2\) para encontrar la clave.
2. Construimos un archivo de claves y IDs de páginas ordenadas, sin datos. La búsqueda binaria en este archivo es
mucho más rápida. Para lidiar con las páginas, podemos construir un árbol de archivos índice, que contiene el primer
índice de la siguiente página en cada nodo:
3. El problema es que la inserción no funciona muy bien. Al insertar vamos agregando páginas al final, hasta que
lentamente las páginas overflow no indexadas terminan siendo la mayoría de la base, y estamos haciendo búsqueda
lineal para encontrarlas.
B+ Trees
1. Nodos similares a ISAM: pares de <Key, PagePtr> con misma invariante. La búsqueda es igual. La diferencia es
que es un arbol de índice dinámico: siempre está balanceado, por lo que los inserts y deletes siempre son eficientes.
2. Insert:
Hacemos búsqueda para encontrar el lugar donde va el registro. Si la hoja donde iría tiene espacio, lo
agregamos ahí y sólo tenemos que ordenar la hoja. Si la hoja está llena, creamos una nueva y dividimos la
anterior en 2; una vez que tenemos el doble de espacio, agregamos el registro nuevo. Luego actualizamos los
pointers next y prev de las páginas(*).
Alt text
Como la página nueva quedó huerfana, repetimos el mismo proceso recursivamente en la fila anterior:
Agregamos un nodo, y distribuimos el anterior en dos. Luego insertamos un nuevo nodo raíz padre con los
valores entre el final del nodo viejo y el inicio del nuevo.
Alt text
(*) Nótese que el B+ Tree, entonces, no es una estructura estática: los nodos no se almacenan de manera
necesariamente secuencial con el tiempo, por lo que cada uno necesita pointers next y prev.
3. Bulk-loading: Primero ordenamos el input por clave, y así podemos llenar cada hoja de a una y actualizar los pointers
del parent. Cuando el parent se llene, hacemos lo mismo que con insert: dividimos el nodo y usamos el pointer del
medio como nuevo parent.
07 Indices II
Factores a considerar en un índice:
Qué clase de queries permite?
Clave de búsqueda elegida
Almacenamiento usado
Trucos de claves de largo variable
Costo de índice vs heap vs archivo ordenado
Selecciones
Básicas: Igualdad, rango. Los B+ trees tienen ambas pero los linear hashes sólo soportan igualdad, aunque igual
tienen sus usos.
Strings: Regex, matches de cadenas de genomas, etc.
Selecciones n-dimensionales: Rectángulos, radios, near-neighbors. R-tree, KD-tree. Ver también los GiST indexes de
Postgres.
Búsqueda y ordenamiento
En un índice ordenado, las claves de organizan lexicográficamente por columna clave de búsqueda. Primero la primera, etc.
Almacenamiento de Data Entries
Datos vs pointers a datos
Almacenamiento en clusters o no respecto del índice
Tres alternativas para data entries:
1. Por valor: Las hojas directamente contienen los datos.
2. Por referencia: Las hojas contienen pointers (E.g. con el formato {key, recordID} que vimos antes, donde recordID es
un pageID y un slotID) a los datos.
3. Por lista de referencias: Lo mismo que 2), pero comprimido en caso de que varios registros compartan keys.
En la práctica, el indexado por referencia es necesario para soportar múltiples índices por tabla. Los indices indexados
pueden estar en clusters también, lo que da beneficios de localidad y soporta mejor ciertos tipos de compresión, a cambio de
ser más dificiles de mantener ordenados.
Costo de Operaciones
Alt text
08 Buffer Management
Adaptador entre el Disk Space Management, que usa el disco, y el Files and Index Management, que opera en RAM. El BM
maneja un pool de buffers en RAM, formado por rangos de memoria del tamaño de un page, llamados frames. El BM mueve
páginas del disco al buffer pool. Si una página pasa al BM y es modificada, pasar a ser lo que llamamos una "Dirty Page".
Funcionamiento del BM
Al bootear, el servidor allocea un rango grande de memoria para el buffer pool. Los frames se mantienen contiguos en esa
memoria, junto con sus metadatos: FrameID, PageID, dirty flag y pin count.
Requests
1- Elegir un frame con pin_count = 0 para reemplazar. 2- Si está dirty, escribir la página actual al disco. 3- Leer la página
solicitada al frame. 4- Pinear la página y volver la dirección.
Luego, el solicitante debe setear el dirty bit si modifica la página y despinear la página lo antes posible. Una página puede
ser solicitada varias veces, por lo que el pin_count se incrementa (Como un Rc<T>). Otras cosas que vamos a ver después
como CC, recovery y Logs pueden usar I/O también.
Page Replacement
Políticas posibles: LRU, Clock. Conviene elegir una según el access pattern a la DB.
LRU: Rastrear tiempo desde el unpin. El frame que se reemplaza es el usado por última vez. Es útil cuando se accede
repetidamente a algunas páginas populares (localidad temporal). La desventaja es que hay que hacer un find min para
encontrar el frame más viejo (generalmente almacenando los datos en un priority queue).
Clock: Heurística aproximada de LRU. Agregamos un bit de referencia a los metadatos, y usamos un iterador cíclico para
recorrer las páginas. Al encontrar una página con el RefBit seteado, lo sacamos y avanzamos el iterador. Al encontrar uno
sin pin ni refbit, reemplazamos esa, seteamos su refbit y avanzamos el clock.
Problemas: LRU/Clock tiene mala performance con scaneos repetidos de archivos grandes. Lo peor que puede ocurrir es
Sequential Flooding del pool, donde la lectura secuencial del archivo hace que siempre se elimine del buffer justo lo que se
estaba por leer. LRU gana en acceso aleatorio, y MRU en acceso secuencial repetido.
Mejoras posibles:
Prefetching: solicitar secuencias de páginas cada vez que se solicita una.
Usar información del DBMS para influenciar el comportamiento del BufMgr: Para queries grandes, el DBMS puede
predecir patrones de I/O de un puñado de algoritmos de procesamiento de queries.
Encontrar mejores algoritmos estocásticos: 2Q, LRU-2, ARC. Ver Page Replacement Algorithm en Wikipedia, o
buscar los papers de estos algos.
Híbridos: Caso especial para indices, LRU-2 por default. PostgreSQL usa Clock.
09 Algebra Relacional
Operadores
Unarios: Proyección (\(\pi\)): Retiene sólo las columnas deseadas (vertical), selección (\(\sigma\)): Subconjunto de
filas (horizontal), Renaming (\(\rho\)): renombrar atributos y relaciones. También se le suele agregar un Group By (\
(\gamma\))
Binarios: Unión (\(\cup\)), diferencia (\(-\)), producto cruz (\(\times\))
Compuestos: Intersección (\(\cap\)) y joins (\(\bowtie_\empty, \bowtie\))
10 Sorting and Hashing
Para qué es el sorting en las queries?
1. Eliminar duplicados, agrupar (DISTINCT / GROUP BY)
2. Ordenamiento explícito (ORDER BY)
3. Primer paso en el bulk-loading de index trees
Algoritmos Out-of-Core
Heurística general: Streamear datos a RAM en una pasada, y usar D&C con fragmentos que entren en RAM.
Double Buffering: Los buffers de Input y Output se cargan y descargan independientemente según se llenen, y un segundo
buffer permite usar el I/O en paralelo mientras el buffer anterior se encuentra en uso. Sirve para cualquier algoritmo de
streaming.
Hashing: En muchos casos no necesitamos orden, sino agrupar registros iguales.
11 Iterators
Podemos pensar cada operador relacional de la query lógica como un iterador:
abstract class iterator {
void setup(List<Iterator> inputs)
void init(args);
tuple next();
void close();
}
El iterador puede ser streaming o blocking, mantener un estado interno, y ser input de cualquier otro iterador.
12 Query Optimization - Plan Space
Optimizadores comúnes: System R, Cascades.
![[Pasted image [Link]]]
Parser: Chequea correctitud y autorización, genera un parse tree.
Rewriter: Convierte queries a forma canónica, achatando views, convirtiendo subqueries en joins, etc. En las DBMS
open-source suelen no ser tan buenos.
"Cost-based" Optimizer: Optimiza queries de a un bloque a la vez (SELECT, JOIN, etc.), usando catalog stats para
encontrar el plan menos costoso por bloque.
Los optimizadores se encargan de tres problemas ortogonales:
Plan space: Qué planes considerar para una query?
Cost estimation: Cómo se estiman sus costos?
Search strategy: Cómo se busca en el espacio de planes?
Equivalencias Algebraicas
Selects: Commutatividad, cascada
Projection: Cascada.
Producto cartesiano: Asociatividad, commutatividad
Joins: Se pueden pensar igual que el producto cartesiano, aunque cuidando el orden.
Heurísticas Comúnes
Selects: Aplicar selecciones en cuanto tengas las columnas relevantes.
Projects: Mantener sólo las columnas necesarias para evaluar los operadores próximos.
Evitar productos cartesianos: Sólo son buenos para unir tablas pequeñas entre sí. En general es mejor hacer theta-
joins.
Equivalencias Físicas
Acceso a tablas únicas: Heap scan, Index scan (si está disponible).
Equijoins: Block Nested Loop, Index Nested Loop, Sort-Merge Join, Grace Hash Join.
No-equijoins: Block Nested Loop.
13 Query Optimization - Costs and Search
Requisitos para optimizar queries:
1. Espacio de planeamiento, basado en equivalencias relacionales e implementaciones distintas de los operadores
2. Estimación de costo, basada en formulas de costo y estimaciones de tamaño, que dependen a su vez de catalogar la
información de las tablas y estimar la selectividad (o Reduction Factor) de las operaciones
3. Un algoritmo de búsqueda para recorrer el espacio de planeamiento y encontrar la opción de menor costo.
Optimizador Salinger (System R)
Funciona bien hasta 10-15 joins
El espacio de planeamiento debe ser recortado: dos heurísticas comunes son evitar productos cartesianos y
considerar sólo planes "left-deep"
Estimación de costo: Inexacta, basada en estadísticas en catálogos del sistema, considerando costos de CPU y de
I/O
Algoritmo de búsqueda: Programación Dinámica
Plan Space
Todas las expresiones algebraicas equivalentes, y las mezclas de implementaciones físicas de esas expresiones
Tenemos que considerar propiedades físicas como sorting
Estimación de Costo
\(\#I/O + CPUfactor \cdot \#tuples\)
Catálogos: Necesitamos info de las relaciones y los índices involucrados. Los catálogos suelen contener al menos
estos valores. Los DBMS modernos suelen incluir más datos, como histogramas de distribución de datos.
Estadística Significado
NTuples # Tuplas en la tabla (cardinalidad)
Npages # de páginas
Low/High Valores Min y Max de la columna
Nkeys # de valores distintos en la columna
IHeight La altura de un índice
INPages # de páginas del disco en un índice
Cardinalidad máxima de output = producto de las cardinalidades input. \(Selectividad = |output|/|input|\). Cada term
tiene su selectividad.
select attributes
from relations
where term1 and .. termk
Cardinalidad resultante = \(MaxNumTuples \cdot ProductoDeSelectividades\)
La selectividad de un predicado de igualdad entre dos columnas es \(1/MAX(NumKeysLeft, NumKeysRight)\)
Ejemplos:
![[Pasted image [Link]]]![[Pasted image [Link]]]
Selectividad de Joins: \(s_p s_q |R| |S|\), donde s y p son términos independientes entre sí para cada tupla.
Algoritmo de Búsqueda
Ejemplo:
![[Pasted image [Link]]]
Programación Dinámica en System R
Ejemplo
SELECT [Link], COUNT(*) AS number
FROM Sailors S JOIN Reserves R JOIN Boats B
WHERE [Link] = "red"
GROUP BY [Link]
Sailors: B+ tree index on sid
Reserves: Clustered B+ tree on bid, B+ on sid
Boats: B+ on color
Pass 1: Mejor plan(es) para cada relación
Sailors, Reserves: File Scan
Interesantes: B+ tree en [Link] y [Link]
Boats: B+ tree en color
Pass 2:
for plan P in pass 1:
for FROM table T not in P:
for access method M on T:
for each join method:
generate P JOIN M(T)
Pass 3+:
Usar planes del pass 2 como relaciones exteriores, generar planes para el siguiente join igual que en pass 2.
Luego, agregar costo de groupby (= costo de ordenar el resultado por sid, salvo que ya haya sido ordenado antes)
Luego, elegir el plan más barato.
14 Parallelism
Métricas de paralelismo:
Speed-up: Aumentar el hardware, workload fijo.
Scale-up: Aumentar el hardware, aumentar el workload.
Tipos de Paralelismo
Pipeline: Escalar hasta la profundidad del pipeline, haciendo que cada estadío trabaje al mismo tiempo con
resultados previos. Limitado por la profundidad del pipeline.
Partition: Escalar una sóla función en cuanta data sea posible al mismo tiempo. Paralelizable tanto como la cantidad
de datos.
Paralelismos de queries
Interquery: Cada query corriendo en un procesador separado. Requiere de control de concurrencias de contemple
paralelismo.
Intraquery:
Interoperador: Pipeline parallelism entre operadores, o "tree" parallelism cuando tenés operadores en ambos
lados de un arbol de joins.
Intraoperador: Partition parallelism para que un sólo operador de join corra en varios nodos.
Joins
One-sided shuffle: si una tabla R ya está particionada, se puede particionar sólo S y hacer joins locales con los R ya
particionados.
Broadcast Join: si R es pequeña, enviar una copia a cada nodo con una partición de S.
Arquitecturas Paralelas
Shared Memory, Shared Disk, Shared Nothing (a.k.a. Cluster).
15 Unstructured Data - Searching Text
Arquitectura
IR DBMS
Imprecise Semantics Precise Semantics
Keyword Search SQL
Unstructured Text Structured Data
Read mostly, add docs in batches Transactional Updates
Page through top k results Generate full answer
![[Pasted image [Link]]]
Modelo "Bag of Words"
Detalle 1: Eliminar stop words (artículos, tags de html, etc.).
Detalle 2: "Stemming". Convertir palabras a forma básica (ver nltk en python).
Almacenamiento
Boolean Text Search: Ej.: "Microsoft" AND ("Glass" OR "Door") AND NOT "Windows"
Text Index: En IR suele tener un significado más abarcativo que en DBs. Generalmente un schema logico (i.e. tablas)
con un schema físico (i.e. indexes, como un postings list), generalmente almacenado en archivos y no en un DBMS.
No suelen tener deletes o modificaciones. Ver "Log-Structured Merge" index (LSM).
Si no se puede desactivar la DB para actualizar el index, se crea un segundo index en otra máquina y se reemplaza
el primero con el segundo; también evita problema de concurrencia.
16 DB Design - ER Models
Ver: ISA hierarchies
17 DB Design - FDs and Normalization
A relationship R is in BCNF if the only non-trivial FDs over R are key constraints.
18-19 Transactions and Concurrency
Transaction
A sequence of multiple actions to be executed as an atomic unit. For the DBMS, it's an abstract view of an activity, a
sequence of reads and writes that must commit or abort as an atomic unit. El Xact manager controla la ejecución de las
transacciones.
ACID
Atomicity: All actions in the Xact happen, or none happen. Goes hand in hand with durability.
Consistency: If the DB starts out consistent, it ends up consistent at the end of Xact. What "consistent" means
depends on the DB, but for a relational DB we can say that consistency is a set of declarative integrity constraints
(primary key, type, FDs).
Isolation: Execution of each Xact is isolated from that of others. Net effect should be equal to executing all
transactions on a sequential order
Durability: If a Xact commits, its effects persist. The effects of a committed Xact must survive failures; this is tipically
done by logging all actions. Aborted transactions can then be undone, and actions of commited transactions that didn't
get to propagate before a crash can be Redone.
Concurrency Control
Naive approach (Serial schedule): Schedule actions of every transactions separately, so that there's effectively no
concurrency.
Two schedules are equivalent if they use the same transactions, with the same ordering, and leave the DB in the
same final state. Ergo, a schedule is serializable if it's equivalent to some serial schedule.
How to check if two schedules leave the DB in the same state? We can check for conflicting operations. The order of non-
conflicting operations has no effect on the final state of the DB. Two operations conflict if they:
Are by different transactions.
Are on the same object.
At least one of them is a write.
In practice, conflict serializability is what gets used, even though it doesn't allow all serializable schedules, because it can be
enforced efficiently. Search "Escrow Transactions" for a system that allows for more concurrency, by understanding some of
the meaning of the data being handled.
Two Phase Locking (2PL)
Esquema "pesimista" para mantener serializabilidad de conflictos. Siempre crea locks por si llega a haber conflicto. Otros
esquemas optimistas permiten avanzar a la transacción y abortan si hubo conflicto. Ver Optimistic CC, Timestamp-Ordered
Multiversion CC.
Reglas de 2PL:
Una Xact necesita un lock compartido (S) antes de leer, y uno exclusivo (X) antes de escribir.
Una Xact no puede adquirir nuevos locks luego de liberar cualquier lock.
Para evitar cascades de aborts, se usa Strict 2PL, que libera todos los locks juntos a la vez cuando se termina la transacción.
Deadlocks
Métodos para evitar deadlocks:
Prevention
Avoidance
Detection and Resolution
Timeouts
Algunos deadlocks son inevitables: Porque hay varias upgrades, o varios locks en transacciones largas. Las técnicas de
ordenamiento de recursos que se usan en sistemas operativos no funcionan para una base de datos compartida; qué orden
se impondría?
Deadlock detection: Mantener un grafo "waits-for" de transacciones, y buscar ciclos en el grafo.
Intent Locks
Granularidad de locks: lockeamos tuplas, páginas, tablas? Cuanto más grande lo que lockeemos, menos locks para manejar,
pero más serializadas van a terminan las transacciones. Hilando más fino quedan más recursos disponibles, pero hay que
manejar muchos locks.
La solución es permitir locks de varios tamaños y definir una jerarquía de granularidades, con los locks pequeños dentro de
los grandes, representado como árbol. Cuando una transacción lockea un nodo del arbol, implicitamente lockea también
todos sus hijos. Introducimos el concepto de "intent" lock: Antes de conseguir un lock S o X, una transacción necesita intent
locks en todos sus ancestros en la jerarquía.
![[Pasted image [Link]]]
Tipos de intent lock:
IS: Intención de lock S en granularidad más fina.
IX: Lo mismo pero de lock X.
SIX: Como S y IX a la vez.
VER: Next key locking. Evitar los "fantasmas" en las queries lockeando rangos lógicos en lugar de sólo datos discretos.
20 Recovery
Steal/No-Force Policy
21 Webcrawlers and IR, Parallel Search and
Ranking
22 Distributed Transactions
23 NoSQL
Cómo escalar una DB:
Partitioning: distribuir la DB entre varias máquinas en el cluster. Esparcimos las queries entre varios servidores y
aumentamos el throughput. Los writes no generan problemas, pero los reads se vuelven más caros.
Replication: Crear múltiples copias de cada partición de la DB. Igual que el anterior y además tiene mejor tolerancia a
las fallas. Los reads son fáciles pero los writes se vuelven más caros. En ambos casos, la consistencia es dificil de
manejar.
El objectivo de NoSQL es simplificar el modelo de datos, cediendo funcionalidad. En lugar de los principios ACID, tenemos
BASE:
Basic Availability: La aplicación debe manejar por sí misma las fallas parciales.
Soft State: El estado de la DB puede cambiar incluso sin inputs.
Eventual Consistency: La DB "eventualmente" será consistente.
Taxonomía de Modelos NoSQL
Key-value stores:
Amazon DynamoDB, Voldemort, Memcached.
Todo se almacena en pares (key, value).
Key = string/int único, Value = cualquier cosa
Sólo dos operaciones: get(key) y put(key, value). Operaciones en el valor no soportadas, debe implementarlas la
aplicación.
Fácil de distribuir sin replicación: Usás un hash para almacenar una key \(k\) en un servidor \(h(k)\). Para replicar,
almacenás una key en \(h1(k), h2(k), etc.\) y en updates propagás los cambios a otros servidores (Eventual
consistency).
Ergo, los key-value stores suelen combinar partitioning y replication.
Ejemplo: Para representar una DB de vuelos, podríamos usar:
1. key = fid, value = registro de un vuelo
2. key = date, value = vuelos de ese día
3. key = (origin, dest), value = todos los vuelos entre origin y dest
Extensible Record / Wide-column Stores:
HBase, Cassandra, PNUTS. Basados en BigTable de Google.
1. key = rowID, value = record
2. key = (rowID, columnID) (o múltiples columnID), value = field.
Document stores:
SimpleDB, CouchDB, MongoDB.
Datos Semiestructurados
Idea: Almacenar valores como datos estructurados (i.e. documentos) para facilitar el procesamiento (e.g. JSON, XML).
Podemos transformar relaciones en documentos (achatándolos de distintas maneras) o documentos en relaciones (más
dificil sin son datos con muy heterogéneos).
CREATE TABLE people (person json);
SELECT * FROM people
WHERE person @> '{"name": "Mary"}';
Los datos semiestructurados son útiles como formato de intercambio. Los sistemas modernos también los usan como
modelo de datos para DBs: SQL Server soporta valores XML, MySQL tiene json y jsonb, además de DBs como CouchBase,
MongoDB y Snowflake que usan JSON, o BigQuery que usa Protobuf.
MongoDB
Permite una evolución rápida del modelo de datos, ideal para startups y web apps con schemas que cambian
frecuentemente.
BASE en lugar de ACID.
No soportaba joins, ahora sólo soporta left outer joins.
Opciones limitadas para optimizar queries.
Actualmente soporta validación de schemas json, aunque no suele usarse.
MongoDB DBMS
Collection Relation
Document Row/Record
Field Column
Document= {field : value, ...}, donde un valor puede ser atómico, un documento, un array de atómicos, o un array de
documentos, igual que el modelo JSON, pero almacenado internamente como BSON (Binary JSON). Las librerías cliente
pueden operar directamente sobre BSON. Además, cada documento tiene un campo especial _id, indexado por default y
agregado automáticamente si no está en la ingesta. MongoDB es una base NoSQL distribuida. Las colecciones se
particionen en base a rangos de un campo. Luego las particiones se replican de manera asincrónica para manejar fallas,
pero los fallos de una partición no propagada se pierden.
MQL
Tipos de queries MQL:
Retrieval: SELECT-WHERE-ETC. restringido
Aggregation: Pipeline general de operadores
Updates
[Link].op1(...).op2(...)...
El input y output de una query MQL son colecciones. Se centra en manipular una sola colección a la vez. Usamos puntos
para bajar en cadenas de docs/arrays, y los campos van entre comillas (e.g. "[Link]" es el campo qty dentro de instock).
También se pueden notar los elementos con números (e.g. "[Link]"). El signo $ indica keywords especiales (e.g. $gt,
$lte, $add, $elemMatch). Se utilizan poniéndolos en el campo de una expresión "field: value", por ejemplo:
{LOperand : { $keyword : ROperand}}
{qty : {$gt : 30}}
// O con arrays en value para varios argumentos
{$add : [1, 2]}
MQL Queries
[Link](<predicate>, <projection>*), con tanto predicate como projection expresados como
documentos. Para proyectar, se especificar campos con un 1 para incluírlos, o con un 0 para excluírlos. Por ejemplo:
find( { status : "D"}, { item: 1, _id : 0 }) // ID es 1 por default
find ( {}, {item : 1, tags : 0 }) // ERROR
find ( {}, {item : 1, "[Link]" : 1, _id : 0 })
find( { $or: [ { status: "D"}, { qty : { $lt : 30 }}]})
Otras funciones:
[Link]({}).limit(500).sort( { "dim.0" : -1 , item : 1 })
Aggregation: Caso general de query, con pipelines.
$group : {
_id: <expression>,
<field1 : { <aggfunc1> : <expression1> },
...
}
[Link]( [
{ $group : { _id: "$state", totalPop: { $sum: "$pop" }}}, // GROUP BY, AGGS, SELECT
{ $match : { totalPop : { $gte : 15000000 }}}, // HAVING
{ $sort : { totalPop : -1 }} // SORT BY
])
Unwind: Expande un array construyendo un documento por elemento del array.
aggregate([
{ $unwind : "$tags" },
{ $project : {_id : 0, instock: 0}}
])
Insert/Update/Delete Many: Insert crea colecciones si están ausentes. Agrega el atributo _id automaticamente y lo pone al
comienzo.
insertMany ( [ { <collection> }, ... ])
updateMany ( { <condition> }, { <change> })
24 MapReduce and Spark
Tipos de paralelismo:
Inter-query: Una query por nodo, bueno para workload transaccionales (OLTP).
Inter-operator: Un operador por nodo, bueno para workloads analíticos (OLAP).
Intra-operator. Operadores en múltiples nodos, bueno para ambos. Algoritmos de join paralelizables: partitioned
Hash-Join, Broadcast Join.
Distributed File System (DFS)
Para archivos muy grandes. Cada archivo se particiona en chunks, generalmente de 64 MB. Cada chunk se replica varias
(>3) veces en racks distintos. Implementado en GFS de Google y HDFS de Hadoop.
MapReduce
Propuesto en 2004, seguido de su implementación open source, Hadoop. Es un modelo de programación e implementación
para procesamiento paralelo a gran escala. Modelo de datos: Archivos. Cada archivo es una bolsa de pares (key, value). Un
programa de MapReduce tiene una bolsa de pares como input, y otra como output. Sirve para solucionar el problema de
lectura de datos a gran escala, de la siguiente manera:
1. Map: Extraer datos que importan de cada registro. El usuario provee una función map. El sistema la aplica en
paralelo a todos los pares (input key, value) del archivo input, resultando en pares (intermediate key, value).
2. En base a los datos de 1), shuffle y sort de los registros.
3. Reduce: "Reducimos" los datos, con algún transformación como agregar, resumir, filtrar, etc. El sistema agrupa todos
los pares con la misma intermediate key, y le pasa los pares (intermediate key, value) a la función reduce, proveída
por el usuario, que devuelve los valores finalmente.
4. Escribir los resultados.
E.g. Contar el número de ocurrencias de una palabra en una colección de documentos. Cada documento tiene key =
documentID y value = set de palabras.
map(key: String, value: String):
for word in value:
emitIntermediate(word, 1)
reduce(key: String, values: Iterator):
sum: int = 0
for v in values:
sum += v
emit(key, sum)
![[Pasted image [Link]]]
Workers
Un worker es un proceo que ejecuta una tarea a la vez. Generalmente hay un worker por procesador, ergo 8 o 16 por nodo.
Para tolerar fallas, los workers mappers escriben sus resultados intermedios al disco. Los reducers leer los archivos como
input; si el servidor falla, los reducers reinician su tarea en otro servidor.
Implementación
Suele haber un nodo lider, que particiona el archivo input en \(M\) splits, por key. El líder luego asigna workers (=servers) a
las \(M\) tareas de mappeo, registrando su progreso. Los workers escriben su output al disco local, particionado en \(R\)
regiones, que conforman \(R\) tareas de Reduce asignadas a los workers. Los workers leer las regiones producidas por el
map y escribir el output al terminar. Un straggler es un worker que tarda demasiado en completar una tarea, ya sea por
problemas de hardware o porque el scheduler le asignó otras tareas. La solución es ejecutar un backup preventivo de las
pocas tareas restantes en progreso.
Implementación de Operadores Relacionales
Selection: No hace falta el reducer; este es uno de los problemas de MapReduce/Hadoop, donde cada map tiene que
tener su reduce.
map(t: Tuple):
if t.A = 123:
emitIntermediate(t.A, t)
reduce(A: string, values: Iterator):
for v in values:
emit(v)
Group By:
map(t: Tuple):
emitIntermediate(t.A, t.B)
reduce(a: String, values: Iterator):
s = 0
for v in values:
s += v
emit(A, s)
Join: \(R(A,B)\bowtie_{B=C}S(C,D)\)
# Partitioned Hash Join
map (t: Tuple):
match [Link]:
'R': emitIntermediate(t.B, t)
'S': emitIntermediate(t.C, t)
reduce(k: string, values: Iterator):
R, S = []
for v in values:
match [Link]:
'R': [Link](v)
'S': [Link](v)
for (v1, v2) in zip(R, S):
emit(v1, v2)
# Broadcast Join
map(value: string): # value = grupo de varias tuplas de R
readFromNetwork(S)
hashTable = new()
for w in S:
[Link](w.C, w)
for v in value:
for w in [Link](v.B):
emit(v, w)
reduce():
pass
Problemas
Dificil de escribir queries complejas. Todo tiene que expresarse como map y reduce.
Las tareas intermedias deben escribir sus resultados al disco por fault tolerance. Muy lento.
Spark
Sistema de procesamiento distribuido sobre HDFS. A diferencia de MapReduce, tiene varios pasos, incluyendo iteraciones;
almacena los resultados intermedios en memoria; y es más cercano al álgebra relacional.
Modelo de Datos
Sus objectos (que pueden ser tuplas, objetos o cualquier cosa) se organizan en RDDs (REsilient Distributed Datasets). Cada
RDD es un dataset distribuido e inmutable, acompañado por su linaje. Un linaje es una expresión que describe cómo se
computó el dataset (e.g. un plan de álgebra relacional). Los resultados intermedios se almacenan como RDD; en caso de un
crash, el RDD se pierde, pero el driver ( = nodo lider) conoce el linaje, por lo que puede recomputar la partición perdida.
Programación en Spark
Un programa de Spark está formado por:
1. transformaciones lazy (map, reduceByKey, join, etc.)
2. acciones eager (count, reduce, save).
Colecciones en Spark:
RDD<T> = Colección RDD de tipo T. Particionada, recuperable (a través del linaje), no nesteada.
Seq<T> = Secuencia. Local a un servidor, puede ser nesteada.
Ejemplo 1: Dado un log hdfs://[Link], recuperar las lineas que comiencen con "ERROR" y contengan el texto
"sqlite".
s: SparkSession = [Link]()...getOrCreate()
lines: JavaRDD<string> = [Link]().textFile("hdfs://[Link]")
errors: JavaRDD<string> = [Link](lambda l: [Link]("ERROR"))
# "Materializador" que hace que Spark compute todo el resultado hasta acá
[Link]()
sqlerrors: JavaRDD<string> = [Link](lambda l: [Link]("sqlite"))
[Link]()
Ejemplo 2: SELECT count(*) FROM R NAT JOIN S WHERE R.B > 200 AND S.C < 100
R = [Link]().textFile("[Link]").map(parseRecord).persist()
S = [Link]().textFile("[Link]").map(parseRecord).persist()
RB = [Link](lambda t: t.b > 200).persist()
SC = [Link](lambda t: t.c < 100).persist()
J = [Link](SC).persist()
[Link]()
Spark 2.0
Agrega DataFrames y DataSets. Los DataSets son como DataFrames, pero en lugar de contener sólo filas, los elementos
deben ser objetos con tipo (DataFrame == Dataset<Row>).
VER What Goes Around Comes Around, Hellerstein & Stonebraker