cassandra and spark, closing the gap between no sql and analytics codemotion berlin 2015

56
@doanduyhai Cassandra & Spark, closing the gap between NoSQL and analytics DuyHai DOAN, Technical Advocate

Upload: duyhai-doan

Post on 24-Jan-2018

1.125 views

Category:

Technology


0 download

TRANSCRIPT

Page 1: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Cassandra & Spark, closing the gap between NoSQL and analytics DuyHai DOAN, Technical Advocate

Page 2: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Who Am I ?Duy Hai DOAN Cassandra technical advocate•  talks, meetups, confs•  open-source devs (Achilles, …)•  OSS Cassandra point of contact

[email protected] ☞ @doanduyhai

2

Page 3: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Datastax•  Founded in April 2010

•  We contribute a lot to Apache Cassandra™

•  400+ customers (25 of the Fortune 100), 400+ employees

•  Headquarter in San Francisco Bay area

•  EU headquarter in London, offices in France and Germany

•  Datastax Enterprise = OSS Cassandra + extra features

3

Page 4: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Spark – Cassandra Use Cases

Load data from various sources

Analytics (join, aggregate, transform, …)

Sanitize, validate, normalize, transform data

Schema migration, Data conversion

4

Page 5: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

Spark & Cassandra Presentation

Spark & its eco-systemCassandra Quick Recap

Page 6: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

What is Apache Spark ?Created at

Apache Project since 2010

General data processing framework

Faster than Hadoop, in memory

One-framework-many-components approach

6

Page 7: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Partitions transformations

map(tuple => (tuple._3, tuple))

Direct transformation

Shuffle (expensive !)

groupByKey()

countByKey()

partition

RDD

Final action

7

Page 8: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Spark eco-system

Local Standalone cluster YARN Mesos

Spark Core Engine (Scala/Java/Python)

Spark Streaming MLLib GraphX Spark SQL

Persistence

Cluster Manager

etc…

8

Page 9: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Spark eco-system

Local Standalone cluster YARN Mesos

Spark Core Engine (Scala/Java/Python)

Spark Streaming MLLib GraphX Spark SQL

Persistence

Cluster Manager

etc…

9

Page 10: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

What is Apache Cassandra?Created at

Apache Project since 2009

Distributed NoSQL database

Eventual consistency

Distributed table abstraction

10

Page 11: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Cassandra data distribution reminderRandom: hash of #partition → token = hash(#p)

Hash: ]-X, X]

X = huge number (264/2)

n1

n2

n3

n4

n5

n6

n7

n8

11

Page 12: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Cassandra token rangesA: ]0, X/8]B: ] X/8, 2X/8]C: ] 2X/8, 3X/8]D: ] 3X/8, 4X/8]E: ] 4X/8, 5X/8]F: ] 5X/8, 6X/8]G: ] 6X/8, 7X/8]H: ] 7X/8, X]

Murmur3 hash function

n1

n2

n3

n4

n5

n6

n7

n8

A

B

C

D

E

F

G

H

12

Page 13: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Linear scalability

n1

n2

n3

n4

n5

n6

n7

n8

A

B

C

D

E

F

G

H

user_id1

user_id2

user_id3

user_id4

user_id5

13

Page 14: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Linear scalability

n1

n2

n3

n4

n5

n6

n7

n8

A

B

C

D

E

F

G

H

user_id1

user_id2

user_id3

user_id4

user_id5

14

Page 15: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Cassandra Query Language (CQL)

INSERT INTO users(login, name, age) VALUES(‘jdoe’, ‘John DOE’, 33);

UPDATE users SET age = 34 WHERE login = ‘jdoe’;

DELETE age FROM users WHERE login = ‘jdoe’;

SELECT age FROM users WHERE login = ‘jdoe’;

15

Page 16: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

Spark & Cassandra Connector

Spark Core APISparkSQL/DataFrame

Spark Streaming

Page 17: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Spark/Cassandra connector architectureAll Cassandra types supported and converted to Scala types

Server side data filtering (SELECT … WHERE …)

Use Java-driver underneath

Scala and Java support. Python support via PySpark (exp.)

17

Page 18: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Connector architecture – Core APICassandra tables exposed as Spark RDDs

Read from and write to Cassandra

Mapping of C* tables and rows to Scala objects•  CassandraRDD and CassandraRow •  Scala case class (object mapper)•  Scala tuples

18

Page 19: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Spark Core https://github.com/doanduyhai/incubator-zeppelin/tree/ApacheBigData

Page 20: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Connector architecture – DataFrame Mapping of Cassandra table to DataFrame

•  CassandraSQLContext à org.apache.spark.sql.SQLContext

•  CassandraSQLRow à org.apache.spark.sql.catalyst.expressions.Row

•  Mapping of Cassandra types to Catalyst types

•  CassandraCatalog à Catalog (used by Catalyst Analyzer)

20

Page 21: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Connector architecture – DataFrame Mapping of Cassandra table to SchemaRDD

•  CassandraSourceRelation •  extends BaseRelation with InsertableRelation with PruntedFilteredScan

•  custom query plan

•  push predicates to CQL for early filtering (if possible)

SELECT * FROM user_emails WHERE login = ‘jdoe’;

21

Page 22: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Spark SQL https://github.com/doanduyhai/incubator-zeppelin/tree/ApacheBigData

Page 23: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Connector architecture – Spark Streaming Streaming data INTO Cassandra table •  trivial setup •  be careful about your Cassandra data model when having an infinite

stream !!!

Streaming data OUT of Cassandra tables (CASSANDRA-8844) ?•  notification system (publish/subscribe)•  at-least-once delivery semantics•  work in progress …

23

Page 24: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Spark Streaming https://github.com/doanduyhai/incubator-zeppelin/tree/ApacheBigData

Page 25: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Q & R

! "

Page 26: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

Spark/Cassandra operations

Cluster deployment & job lifecycleData locality

Page 27: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Cluster deployment

C* SparkM SparkW

C* SparkW

C* SparkW

C* SparkW

C* SparkW

Stand-alone cluster

27

Page 28: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Cassandra – Spark placement

Spark Worker Spark Worker Spark Worker Spark Worker

1 Cassandra process ⟷ 1 Spark worker

C* C* C* C*

28

Spark Master

Page 29: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Cassandra – Spark job lifecycle

Spark Worker Spark Worker Spark Worker Spark Worker

C* C* C* C*

29

Spark Master

Spark Client

Driver Program

Spark Context 1

D e f i n e y o u r business logic here !

Page 30: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Cassandra – Spark job lifecycle

Spark Worker Spark Worker Spark Worker Spark Worker

C* C* C* C*

30

Spark Master

Spark Client

Driver Program

Spark Context

2

Page 31: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Cassandra – Spark job lifecycle

Spark Worker Spark Worker Spark Worker Spark Worker

C* C* C* C*

31

Spark Master

Spark Client

Driver Program

Spark Context

3 3 3 3

Page 32: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Cassandra – Spark job lifecycle

Spark Worker Spark Worker Spark Worker Spark Worker

C* C* C* C*

32

Spark Master

Spark Client

Driver Program

Spark Context

Spark Executor Spark Executor Spark Executor Spark Executor

4 4 4 4

Page 33: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Cassandra – Spark job lifecycle

Spark Worker Spark Worker Spark Worker Spark Worker

C* C* C* C*

33

Spark Master

Spark Client

Driver Program

Spark Context

Spark Executor Spark Executor Spark Executor Spark Executor

5 5 5 5

Page 34: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Data Locality – Cassandra token rangesA: ]0, X/8]B: ] X/8, 2X/8]C: ] 2X/8, 3X/8]D: ] 3X/8, 4X/8]E: ] 4X/8, 5X/8]F: ] 5X/8, 6X/8]G: ] 6X/8, 7X/8]H: ] 7X/8, X]

n1

n2

n3

n4

n5

n6

n7

n8

A

B

C

D

E

F

G

H

34

Page 35: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Data Locality – How ToC*

SparkM SparkW

C* SparkW

C* SparkW

C* SparkW

C* SparkW

Spark partition RDD

Cassandra tokens ranges

35

Page 36: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Data Locality – How ToC*

SparkM SparkW

C* SparkW

C* SparkW

C* SparkW

C* SparkW

Use Murmur3Partitioner

36

Page 37: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Read data localityRead from Cassandra

37

Page 38: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Read data localitySpark shuffle operations

38

Page 39: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Write to Cassandra without data locality

Async batches fan-out writes to Cassandra

Because of shuffle, original data locality is lost

39

Page 40: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Write to Cassandra with data locality

Write to Cassandra

rdd.repartitionByCassandraReplica("keyspace","table")

40

Page 41: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Write data locality•  either stream data in Spark layer using repartitionByCassandraReplica() •  or flush data to Cassandra by async batches•  in any case, there will be data movement on network (sorry no magic)

41

Page 42: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Joins with data localityCREATE TABLE artists(name text, style text, … PRIMARY KEY(name));

CREATE TABLE albums(title text, artist text, year int,… PRIMARY KEY(title));

val join: CassandraJoinRDD[(String,Int), (String,String)] = sc.cassandraTable[(String,Int)](KEYSPACE, ALBUMS) // Select only useful columns for join and processing .select("artist","year") .as((_:String, _:Int)) // Repartition RDDs by "artists" PK, which is "name" .repartitionByCassandraReplica(KEYSPACE, ARTISTS) // Join with "artists" table, selecting only "name" and "country" columns .joinWithCassandraTable[(String,String)](KEYSPACE, ARTISTS, SomeColumns("name","country")) .on(SomeColumns("name"))

42

Page 43: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Joins pipeline with data locality

LOCAL READ FROM CASSANDRA

43

Page 44: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Joins pipeline with data locality

REPARTITION TO MAP CASSANDRA REPLICA

44

Page 45: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Joins pipeline with data locality

JOIN WITH DATA LOCALITY

45

Page 46: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Perfect data locality scenario•  read localy from Cassandra •  use operations that do not require shuffle in Spark (map, filter, …)•  repartitionbyCassandraReplica()•  à to a table having same partition key as original table

•  save back into this Cassandra table

Sanitize, validate, normalize, transform data USE CASE

46

Page 47: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

Spark/Cassandra use-case demos

Data cleaningSchema migration

Analytics

Page 48: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Use Cases

Load data from various sources

Analytics (join, aggregate, transform, …)

Sanitize, validate, normalize, transform data

Schema migration, Data conversion

48

Page 49: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Data cleaning use-casesBug in your application ?

Dirty input data ?

☞ Spark job to clean it up! (perfect data locality)

Sanitize, validate, normalize, transform data

49

Page 50: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Data Cleaning https://github.com/doanduyhai/incubator-zeppelin/tree/ApacheBigData

Page 51: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Schema migration use-casesBusiness requirements change with time ?

Current data model no longer relevant ?

☞ Spark job to migrate data !

Schema migration, Data conversion

51

Page 52: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Data Migration https://github.com/doanduyhai/incubator-zeppelin/tree/ApacheBigData

Page 53: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Analytics use-casesGiven existing tables of performers and albums, I want to :

•  count the number of albums releases by decade (70’s, 80’s, 90’s, …)

☞ Spark job to compute analytics !

Analytics (join, aggregate, transform, …)

53

Page 54: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Analytics pipeline①  Read from production transactional tables

②  Perform aggregation with Spark

③  Save back data into dedicated tables for fast visualization

④  Repeat step ①

54

Page 55: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Analytics https://github.com/doanduyhai/incubator-zeppelin/tree/ApacheBigData

Page 56: Cassandra and Spark, closing the gap between no sql and analytics   codemotion berlin 2015

@doanduyhai

Thank You @doanduyhai

[email protected]

https://academy.datastax.com/