The Deconstructed Database and the Advent of the Open Data Lake
Julien Le Dem: Principal Engineer at Datadog
2
Julien Le Dem�julien.ledem.net
sympathetic.ink
Principal Engineer at
Technical Advisory Council
About the Speaker
3
Agenda
1 The Origin of Times
2 Meanwhile: The Evolution of the Ecosystem
3 Towards Composable Data Systems
The Origin of Times
1
At the Beginning There was Hadoop
Execution: �Map/Reduce
Storage: �Distributed File System
Moving Compute to Storage
M
M
M
R
R
R
Read locally
Write locally
Shuffle
Hadoop Map/Reduce
Great at looking for a needle in a haystack
Hadoop Map/Reduce
… with snow plows
“MapReduce: A major step backwards”
Databases Have Been Around a Long Time
SQL
Declarative
Standard
Consistency constraints
Schema Evolution
Query Evaluation
Syntax
Semantic
Optimization
Execution
Query Evaluation
Syntax
Semantic
Optimization
Execution
SELECT f.c, AVG(b.d)
FROM FOO f
JOIN BAR b ON f.a = b.b
GROUP BY f.c
WHERE f.d = x
Catalog
SELECT
SCAN foo
JOIN
SCAN bar
FILTER
GROUP BY
SELECT
SCAN foo
JOIN
SCAN bar
FILTER
GROUP BY
RDBMS�vs
Map/Reduce
Database
Complex
Somewhat inflexible�Vertically-integrated stack.
�
Map/Reduce
Simple (Simplistic?)�Flexible
Composable
RDBMS
Map/Reduce
You can Even Implement SQL With M/R
Syntax
Execution
SELECT f.c, AVG(b.d)
FROM FOO f
JOIN BAR b ON f.a = b.b
GROUP BY f.c
WHERE f.d = x
HMS
SELECT
SCAN foo
JOIN
SCAN bar
FILTER
GROUP BY
SCAN bar
SCAN foo
SELECT
GROUP BY
JOIN
FILTER
We Used to Build Databases From Scratch
SQL
Initial Release of Hive
Parse query
Transform to a chain of M/R jobs
Done!
Store in the Hadoop File System
It works but …�What does it take to make it efficient?
17
Optimizer
Produce the Best plan
Take advantage of push downs
Interchange
Fast and efficient evaluation
Interchange
Parallelize execution
Fast data exchange
Vectorization
Efficient storage layer
Quick access for OLAP queries
Optimizer
Vectorization
Interchange
Column Store
Optimizer:
Two levels
Rule based:
Things you always want to do.
Ex: Filter before join
Cost based:
Things that depend on context
Ex: Join ordering
Rule Based
Cost Based
Vectorized Execution:
Moving From Row Oriented to Columnar
MonetDB/X100
https://www.cidrdb.org/cidr2005/papers/P19.pdf
Processor Pipeline
Processor Cache
CPU
Cache
Core
Cache
~100 CPU cycles
Core
Efficient Data Exchange
Columnar Storage
Over the Next Decade we did Just That
Parse query
Cost Optimize plan
Resource pools
Columnar in-memory
Code generation
Columnar Storage
Streaming engine
Caching
Join algorithms
Push downs
Done?
The Evolution of the Ecosystem
2
Limitations of Hadoop
Immutable data, �Lack of transaction
Tying storage�and compute
Reliant on discrete steps persisted on disk
Data is Immutable
Turns out Transactions are Useful.
Who knew?
Turns out Transactions are Useful.
He knew
We’re Coupling Storage and Compute
Meanwhile…
31
Technology has Evolved
Networks
The network gets faster faster than the disks. The ratio changed.
No need to move compute to storage.
Computers
More cores
More memory
More disks
More parallelism
Networks
The network got faster, faster than the disks. The ratio changed.
Moving compute to storage became unnecessary
Networks
Computers
The Cloud has Become Ubiquitous
Networks
The network gets faster faster than the disks. The ratio changed.
No need to move compute to storage.
Computers
Networks
On demand
Decoupled Storage
Decoupling Storage and Compute
New School
Blob storage
Old School
Storage and compute tied together
Node 1
Comp
Stor.
Node 2
Comp
Stor.
Node 3
Comp
Stor.
Node 1
Comp
Stor.
Node 2
Comp
Stor.
Node 3
Comp
Stor.
Fast networking
Emergence of Cloud-First
Data Warehouses
Networks
The network gets faster faster than the disks. The ratio changed.
No need to move compute to storage.
Computers
Proprietary format:
Decouple storage and compute:
Great!!!
But…
😇
😈
New Silos
Analysts
Data Scientists
ML Engineers
SQL
Data Frame
Training
Image credit: Justicon
ETL
37
Towards Composable Data Systems
3
What are the Trade Offs?
39
Volume
Latency
Precision
Rebuilding the Database From Components
Governance / Lineage
Where is data going?
Where is it coming from?
Can I guarantee quality?
Calcite:
A Database with BYO Internals
Standard Plan Representation
Transpiling�language independence
Facilitates�Push downs
Reduces complexity�Simplifies integrations
Standard In-Memory Representation
Vectorized execution
Zero-copy network exchange
Cross-language compatible
Vectorization
Takes advantage of modern CPUs
Embeddable
Extensible
You can focus on what makes your application special
The Storage Layer:
A Distributed Dataset With Push Downs
Distributed read:
Push downs:
Storage layer:
Two Levels of Abstraction With Different Trade-Offs
Arrow Flight:
Iceberg & Parquet:
Arrow Flight:
Build Your own Storage Service Layer
Flight Service
Execution engine
Node 1
Node 2
Node 3
Node 4
Node 1
Node 2
Node 3
Node 4
The Open Data Lake
Own your data
Proper Table abstraction
Scalable, Cheap
BLOB Storage
Iceberg: A Proper Table Abstraction
50
Snapshot Isolation:
Partitioning abstraction:
Efficient push downs:
Iceberg:
The Service *is* Your Blob Store
Execution engine
Node 1
Node 2
Node 3
Node 4
Column chunks
Iceberg metadata
The Open Data Lake
Centralized
Specialized
The Open Data Lake
Conclusion and Looking Forward
Libraries enable extensibility
On Demand resources
Formats enable composability
Open Source
Cloud Services
The Future of Cloud Storage?
Bulk updates vs Streaming ingest
Streaming starts adopting Blob storage
Caching/Materialization
Duality of query cache and materialized views
Tables vs Blobs
Blob storage is starting to be table aware
Tables vs Blobs
Bulk updates vs Streaming ingest
Caching / Materialization
Standalone Caching Service?
Cache Service
Pre-agg
Optimize for queries (columnar,sort, bucketing …)
Execution engine
Node 1
Node 2
Node 3
Node 4
Blob Table Storage With Push Downs and Caching?
Flight Service
Execution engine
Node 1
Node 2
Node 3
Node 4
Blob Storage With Streaming Ingest?
Kafka comp. Service
Flight Service
Ingest
Push downs
Optimize for queries
(columnar conversion, sort, bucketing …)
Raw ingest format
Thanks :)