1 of 59

The Deconstructed Database and the Advent of the Open Data Lake

Julien Le Dem: Principal Engineer at Datadog

2 of 59

2

Julien Le Dem�julien.ledem.net

sympathetic.ink

Principal Engineer at

  • Officer of the ASF
  • Member of the LFAI&Data

Technical Advisory Council

  • Co-creator of Parquet, Arrow, OpenLineage
  • Contributed to Calcite, Drill, Iceberg, …
  • Co-Founded a data observability company
  • First Used Hadoop at Yahoo! In 2007

About the Speaker

3 of 59

3

Agenda

1 The Origin of Times

2 Meanwhile: The Evolution of the Ecosystem

3 Towards Composable Data Systems

4 of 59

The Origin of Times

1

5 of 59

At the Beginning There was Hadoop

Execution: �Map/Reduce

Storage: �Distributed File System

6 of 59

Moving Compute to Storage

M

M

M

R

R

R

Read locally

Write locally

Shuffle

7 of 59

Hadoop Map/Reduce

Great at looking for a needle in a haystack

8 of 59

Hadoop Map/Reduce

… with snow plows

9 of 59

“MapReduce: A major step backwards”

  1. A giant step backward in the programming paradigm for large-scale data intensive applications�
  2. A sub-optimal implementation, in that it uses brute force instead of indexing�
  3. Not novel at all -- it represents a specific implementation of well known techniques developed nearly 25 years ago�
  4. Missing most of the features that are routinely included in current DBMS�
  5. Incompatible with all of the tools DBMS users have come to depend on

10 of 59

Databases Have Been Around a Long Time

SQL

Declarative

Standard

Consistency constraints

Schema Evolution

11 of 59

Query Evaluation

Syntax

Semantic

Optimization

Execution

12 of 59

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

13 of 59

RDBMS�vs

Map/Reduce

Database

Complex

Somewhat inflexible�Vertically-integrated stack.

Map/Reduce

Simple (Simplistic?)�Flexible

Composable

RDBMS

Map/Reduce

14 of 59

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

15 of 59

We Used to Build Databases From Scratch

SQL

16 of 59

Initial Release of Hive

Parse query

Transform to a chain of M/R jobs

Done!

Store in the Hadoop File System

17 of 59

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

18 of 59

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

19 of 59

Vectorized Execution:

Moving From Row Oriented to Columnar

MonetDB/X100

https://www.cidrdb.org/cidr2005/papers/P19.pdf

20 of 59

Processor Pipeline

21 of 59

Processor Cache

CPU

Cache

Core

Cache

~100 CPU cycles

Core

22 of 59

Efficient Data Exchange

23 of 59

Columnar Storage

24 of 59

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?

25 of 59

The Evolution of the Ecosystem

2

26 of 59

Limitations of Hadoop

Immutable data, �Lack of transaction

Tying storage�and compute

Reliant on discrete steps persisted on disk

27 of 59

Data is Immutable

28 of 59

Turns out Transactions are Useful.

Who knew?

29 of 59

Turns out Transactions are Useful.

He knew

30 of 59

We’re Coupling Storage and Compute

31 of 59

Meanwhile…

31

32 of 59

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

33 of 59

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

  • Fast
  • Cheap
  • Decoupled from compute

Networks

  • Elastic resources
  • Pay per use

On demand

Decoupled Storage

34 of 59

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

35 of 59

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:

  • Data only accessible from within
  • import/export costs
  • Creation of Silos

Decouple storage and compute:

  • Pay per use
  • Grow on demand
  • Little resource contention

Great!!!

But…

😇

😈

36 of 59

New Silos

Analysts

Data Scientists

ML Engineers

SQL

Data Frame

Training

Image credit: Justicon

ETL

37 of 59

37

38 of 59

Towards Composable Data Systems

3

39 of 59

What are the Trade Offs?

39

Volume

  • Does it fit on one machine?
  • Does it fit in memory or on disk?

Latency

  • Real-time
  • Bulk ingestion

Precision

  • Do approximate results suffice?
  • Is sampling appropriate?

40 of 59

Rebuilding the Database From Components

41 of 59

Governance / Lineage

Where is data going?

Where is it coming from?

Can I guarantee quality?

42 of 59

Calcite:

A Database with BYO Internals

43 of 59

Standard Plan Representation

Transpiling�language independence

Facilitates�Push downs

Reduces complexity�Simplifies integrations

44 of 59

Standard In-Memory Representation

Vectorized execution

Zero-copy network exchange

Cross-language compatible

45 of 59

Vectorization

Takes advantage of modern CPUs

Embeddable

Extensible

You can focus on what makes your application special

46 of 59

The Storage Layer:

A Distributed Dataset With Push Downs

Distributed read:

  • Efficiently format for vectorization
  • Parallel read from multiple streams

Push downs:

  • Projections
  • Filters
  • Limits
  • Aggregations

47 of 59

Storage layer:

Two Levels of Abstraction With Different Trade-Offs

Arrow Flight:

  • Enables finer grain implementation
  • Build your own storage: low latency, transactional, …
  • Flexibility at a cost

Iceberg & Parquet:

  • The standard Open Data Lake
  • Great for bulk updated, scalable & cheap storage
  • No need to operate your service

48 of 59

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

49 of 59

The Open Data Lake

Own your data

Proper Table abstraction

Scalable, Cheap

BLOB Storage

50 of 59

Iceberg: A Proper Table Abstraction

50

Snapshot Isolation:

  • No phantom read
  • Consistency
  • Roll back

Partitioning abstraction:

  • Re-partition as needed
  • Not coupling incremental processing and partitioning

Efficient push downs:

  • Projection
  • Filter
  • Limit

51 of 59

Iceberg:

The Service *is* Your Blob Store

Execution engine

Node 1

Node 2

Node 3

Node 4

Column chunks

Iceberg metadata

52 of 59

The Open Data Lake

Centralized

  • Data
  • Catalog
  • Governance
  • Lineage

Specialized

  • Compute engine

53 of 59

The Open Data Lake

54 of 59

Conclusion and Looking Forward

Libraries enable extensibility

  • Calcite
  • Arrow kernels
  • DataFusion

On Demand resources

  • Storage: reliable, cheap, tiered
  • Compute: elastic

Formats enable composability

  • Substrait
  • Arrow
  • Parquet
  • Iceberg
  • OpenLineage

Open Source

Cloud Services

55 of 59

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

56 of 59

Standalone Caching Service?

Cache Service

Pre-agg

Optimize for queries (columnar,sort, bucketing …)

Execution engine

Node 1

Node 2

Node 3

Node 4

57 of 59

Blob Table Storage With Push Downs and Caching?

Flight Service

Execution engine

Node 1

Node 2

Node 3

Node 4

58 of 59

Blob Storage With Streaming Ingest?

Kafka comp. Service

Flight Service

Ingest

Push downs

Optimize for queries

(columnar conversion, sort, bucketing …)

Raw ingest format

59 of 59

Thanks :)