data.bythebay.io: Monte Zweben, On the path to data Nirvana: Supporting OLTP and OLAP on Hadoop
what um what I'm here to talk about is a new database and it's a database that can support data intensive applications that have a mixed workload both analytical processing as well as powering applications with operational um workloads and the value proposition we provide to our customers is that typically across very many verticals it's very difficult for companies to keep the decisions that they're making in the moment meaning the data that they're using to derive any decisions or analytics is usually old data it's coming from um days ago in fact we ask people all the time in surveys how old is the data in your in your reports simple question and and you'd be surprised what the answer is 75% of the people say it's a day or older even in this day of streaming applications lots of talk about kfka streaming and and Spark streaming and Flink today it's still old data is being used to either learn on or or to perform any kind of analytics and the reason for that is very straightforward that in the old architectures in databases we bifurcated things into two basic systems we bifurcated things into transactional or operational systems powered by traditional databases on the high-end like Oracle And On The Low End like postgress and mice SQL um and then we extract transform and load through the ETL pipelines over to some other analytical framework and that typically is at the high-end terod data or n teaser and and perhaps um some other kind of columnar database over on the olap side but the problem with this is that these pipelines are incredibly complicated these pipelines take a long time and they typically end up being a Big Mo between what's being analyzed and what's being um what's being executed in the business in a real-time fashion and and we think that we're finally at a point where you can reach nirvana meaning that you don't need to have a bifurcated environment you can actually perform your operational workloads whether that's a mobile mile app or web app or um some kind of social app or something that might um be even more advanced and on that same data platform be analyzing in parallel now the old school of oltp and oap is prominent throughout Enterprises but many of you if you've built applications in the last few years that are data intensive machine learning applications you've probably implemented something like a Lambda architecture you guys know what a Lambda architecture is no so so some of you do so a Lambda architecture is essentially um an architecture that takes the real time data coming in and streams it or transfers it to multiple data engines at the same time one data engine that can perform big batch analytics on it typically something that might be on Hadoop using Hive or spark or something like that and then another layer that might be a speed layer where you can perform very quick operations and cleansing operations on the data maybe look up that data and then of course a serving layer that powers an application typically people have used no SQL databases like Cassandra or hbase and there might even be a traditional relational database off on the side powering the application and you can look at this diagram and say why would anyone want to do this when you just look at the architectural comp complexity of keeping all these databases in sync and there are benefits to doing that in the sense that you get access to your data that's real time to power the application perform business logic and at the same time in the background be constantly learning or performing arations and and doing analytics but this is very complex to maintain and so um essentially if you could take all of those data engines and combine them into one data engine that is taking a streaming input that would be Nirvana for the new Lambda architectures and so um these have fewer components it's easier to synchronize it's easier to maintain and you have more application um capability so I'm going to make the case that there is a new kind of architecture that you can utilize for both um combining LP and oap on one data engine or creating what I'll call the Lambda r architecture for a relational Lambda architecture and it's it's all being um enabled by essentially two big technology Trends first it's the scaleout technologies that we see on Hadoop and Spark and of course it's the inmemory technology that we see with systems like spark and um netting it all out the system that I'm here talking about is the splice machine relational database management system and people don't really talk about rdbms's anymore because they kind of went out of style when scale out no SQL came in but they threw the baby out with the bathat when no SQL emerged and got popular because SQL has tremendous power for the data scientist for the application developer it it provides a set of capabilities that for 25 30 years people have leveraged for building applications and those capabilities are still required instead of it being in the database layer we see it ending up in the application layer in spaghetti logic and that's um not the best way to architect systems so we built what we think is the best of both worlds we retain SQL full ANC SQL with asset transactions window functions and all of the capabilities that you would expect in a true relational database but we did it on a scaleout architecture so what that means is you don't store all your data on one system you Shard the data across many different commodity servers it auto shards the data so that you don't have to think about how to split up the data and it's elastic so that you can add more servers as your data loads as your workloads increase and it's incredibly price performant meaning at as you add servers these servers are inexpensive commodity servers you get the power of the parallelism of having all of these servers together but without the cost of heavy engineered systems and this does through parallelism give you 10x speed improvements or better and it also provides you the ability to in a very efficient way and I'll show you how that works perform the mixed workload of being able to perform operational short reads and writes that a real application needs in a concurrent fashion keeping the database consistent and at the same time being able to run analytics aggregations sorts and all kinds of of things that you would expect to do in a more analytical framework or on a big Hadoop based system so what are the kinds of use cases that people use a database like this for well first we see a lot of omni Channel marketing applications what does that mean you have a 360° view of your customer you want to keep track of every interaction you have with your customer whether that's a traditional set of orders that they're buying whether online or in a store um or if that's um interactions in social media or interactions with your mobile web application or perhaps uh even just traditional calls in a call center keeping that into one unified um customer profile and then being able ble to run campaigns that are both sort of batch driven where every so often at some particular marketing drum beat you're reaching out to that customer uh audience either through social media or email or trigger-based Communications where as soon as a a customer does something like purchases a product and gets them within a particular tier of a loyalty program you'd love to be able to communicate directly with them and say hey just 50 more points and you get all of these benefits in that very real time in the- moment way that's Omni Channel marketing and there are lots of systems that do Omni Channel marketing and one particular system I show here up on the slide is an IBM system called Unica it's a third party application unuka is a system that was built by a startup company bought by IBM and um and typically customers use Uno with ETL systems like abono reporting systems like cognos and um we've got customers who actually um compared what it would be like to run this application on Oracle versus running it on spice machine and we saw 10 times um price performance improvements on early versions of of the system and we think um even better with new versions of the system at delivering this on a quarter of the of the um cost so it's about three to seven times um faster depending on the kinds of queries and this was directly against a pretty heavily engineered Oracle rack system so what what the takeaway is here you can use Hadoop and scale out architectures to replace traditional database applications that are running both um large analytical workloads and needing concurrent transactions um we have other customers that are doing slightly different things one customer is a financial um Services um Network broker dealer Network called satera they needed to also have a unified customer profile but in the financial services World a little bit different having a single source of truth that Aggregates many different systems that have asset positions for those clients and being able to calculate advisor commissions based on trades that are happening and being able to do that with slas that get these reports out to the right Executives at the right time was very difficult for them beforehand with all of these disparate systems and they had a very large ETL Pipeline and operating um um operational reporting mechanism and now that's all being powered by the spice machine rdbms and then lastly um showing you another kind of application which is a very modern application a a machine learning application um where um a database of clinical trials is kept which has very deep biomarkers in it across many different types of information um ranging from Clinical information to um genomic information and having that be trained in a set of machine learning models decision Trees Deep learning models and then using those models to take new patients and match them against the clinical trials to determine which clinical trials are best suited for that particular patient that's amazing because what that does is it makes any particular doctor aware of the in the drug Innovations that's going on around the world this could save lives this goes live in July so lots of different use cases across marketing Financial Services Life Sciences the unique aspect of an engine like this is that unlike traditional relational databases where when you're running big analytical loads even just big reports your transactions tend to get backed up behind that analytical workload and you often find yourself even in the best operated scenarios wondering why are the screen getting so slow why can't I pull up this particular record and it's because there's a big analytical report that's running with splice machine we've done something unique we've created a dual engine like the Lambda architecture that had the developer having to synchronize different databases that do different things well we took that under the covers in our database and we have one lane for transactions where transactions can just go take their own computational resources and run and that takes place on a on a data engine called hbase some of you may be familiar with hbas as the key value store on top of Hado and then when you have large joins aggregations groupings and other analytical processes they get run on spark under the covers and I'll show you how that works in a in a minute or two but the whole idea here is that by having workload isolation your operational workloads can operate quickly without being backed up by the analytical workloads which means that as you increase your analytical workload you won't see an exponential increase in the response times of your transactions it should stay fairly constant so that is the hope of and the dream of having resource isolation and being able to do simultaneous oltp and oap workloads and this is important for other new use cases that are emerging many of you have probably heard about the um large um data volumes that are produced by Internet of Things applications other streaming applications this is the kind of architecture that can be connected to Kafka event qes and and uh flow through streaming systems like spark streaming and be able to directly run queries against streams um on a continuous basis because it can ingest data separate from the business logic you can actually perform pretty amazing uh new kinds of applications that are streaming as well so how is this thing built um we relied on three open-source libraries uh we relied on Apache Derby to provide us with the bones of a SQL processor and we gutted Derby's storage layers and planning mechanisms to put in hbas as the uh persistent storage and then Incorporated many different distributed algorithms to enable us to take advantage of hb's distributed architecture and most recently we've Incorporated Apache spark to execute the olap types of queries and I'll show you in a few minutes how it works but the basic architecture is very straightforward basically uh you take uh open a connection via obbc or jdbc to our database like you would with Oracle or postgress um we create an abstract syntax tree of the SQL we create an access plan for that we optimize that access plan using cost based heuristics and advanced approximation um algorithm sketching algorithms like hyper hyper log and then what we do there is pick the right access path the right join order the right join algorithm and the indexes that are necessary and then we execute the queries either on hbas or in spark depending on um what the proper execution path is um some of you are probably familiar with hbas but just to give you a quick little summary hbas is built um to uh reflect the same architecture as Google's big table it is an LSM tree architecture meaning it is very much right optimized it can ingest huge amounts of data data comes directly into memory it builds a memory store it puts the um the uh data on the proper server and the servers are sharded where there is a region of data that is sorted by a primary key on each of these region servers and then um when the memory is uh full then what hbas does is flush the memory out to an indexed Hadoop file and an immutable indexed hute file that later as more and more files are created are compacted into one efficient um immutable file it's a very very nice architecture for doing very short reads and writes in a completely scalable um environment and that's what we use for our main storage many of you are probably very aware of the recent um popularity of Apache spark sparc is providing unmatched performance in uh many different computational um um computational areas machine learning and others um it is an inmemory system it can spill to dis when data sets are too large so it doesn't fail if you can't fit in memory it is very resilient to node failures unlike Hadoop hbase that replicates to provide resiliency um spark does it with ancestry models and it can recalculate Uh computations u based on knowing the lineage of a computation and it's probably the most active Apache Community um right now and it has an incredibly nice set of libraries around it that utilize this technology and so the way we architect is very straightforward um every physical node in a splice machine in implementation has an hbas region server and a spark worker and of course a data node underneath it for hdfs and what happens is is that we can automatically create um spark rdds directly from hbas very very efficiently and we use the same bite code that gets created U for a SQL execution plan in each of these different data flow engines so here's um how it basically works as you um install you basically are getting as I mentioned a region server and a spark worker on each architecture um we distribute basically the um executors out to each of the servers and we don't use any map reduce in order to do the distributed computation it either takes place using spark uh Transformations or it takes place using hbase co-processor so it's very efficient so um the way it works is that you take an a SQL statement you parse the SE the SQL statement create an AB abstract syntax tree we optimize the plan like I mentioned earlier pushing predicates all the way down close to the um close to the um to the um data on the FES and then we unroll nested queries to make it very efficient and then um depending on um what the execution estimates are in the cost-based optimizer if the result sets are very small meaning you're doing a small read and write or maybe a small range scan the bite code that comes out of the optimizer and gets distributed to hbas through co-processors and just execute it in parallel on H bases servers if it's a large result set that is expected perhaps like it's a 20 way join it's aggregating with functions and it's grouping and sorting um that will get sent to spark workers and they will um break that up into batches and use its Fair scheduler to optimize and parallelize that process so it's very clean it's under the covers and it allows you to use the right engine for the right computation and as I mentioned earlier we do support pretty granular Resource Management so using cgroups in Linux you can allocate the proper CPUs memory and priorities to the spark process and to the hbase process so if you have um a circumstance where you really want your transactions to be completely unaffected um by any analytics you might um allocate your memory and CPUs according accordingly additionally spark has a fair scheduler as I mentioned a moment ago and um this is very flexible and customizable so for example we've set some Fair schuer pools these are categories of priorities that you can assign queries to you can assign user groups and roles to these Fair scheduling groups and essentially this the the Sparks Fair scheduler will allocate work to spark workers according to these rules and these priorities and you can extend that declaratively with XML in the setup files and through the um interfaces that uh are associated with spark that give you that really granular capability of controlling how the database performs based on the different queries that are submitted by the different people in the different organizations and one of the other important important things to be able to do is really see what a database is doing especially a a vast distributed database trying to understand you know how it's breaking up tasks and distributing it to different platforms different servers is a difficult thing typically and what we have is a query management UI that lets you see not only what's running now but what's run previously historically you can see a timeline here showing um different queries you can click on them and see the actual query plan and you can see the stages of a query execution and it visualizes the parallelism so you can see which part of a query is being executed in parallel and where there's a barrier where two pieces of parallel operations need to come together and be joined together by some um function like a join and you can dig in really deep into each state AG of these query plans to see you know individual tasks and actually look at granular resources like deserialization time for getting something off of The Wire how much garbage collection is going on how how much how many records are being either read or written to the data dur during this particular task it allows you to get a very granular look at what the database is doing and with these very large data volumes this is very helpful to be able to tune your your database perhaps give hints to plans or actually maybe create new indexes and and key uh your database in perhaps in different ways and provides insight to the design um for the dbas and the data Architects so another really nice feature is in this spark in Hadoop world there's lots of stuff going on and um being able to co exist with all of the different libraries and computational engines that are going that are being invented and being used is really important and what we've done is created a Federated query support so what that means is imagine you wanted to write a query that joins something in splice machine against something that might be in another data format out on Hadoop or maybe even with Oracle um you can do that very declaratively down here in the bottom you can see a little snippet of code where you can see um a select statement where you're just declaring what the data types are in the external U virtual table that you're connecting to and that's all you really need to do in order to bring data in from other sources and we already have you know access to you know for example Amazon files or Oracle and other um other libraries out there we also fit nicely in the um world of spark and map ruce meaning um perhaps somebody wants to run a large machine learning analytic um the features for that machine learning analytic may be stored inside of a spice uh machine relational database um per they may want to perform a little cleansing or a statistical routine on that they can actually run this distributed job um and do so accessing data consistently with transactional Integrity out of spice machine because we provide map input for output format classes so we look like just a normal map reduce job or um we declare to age catalog and you can run Hive against us and you can run spark jobs it's really clean to be able to do that and like I said um the database is really ansy SQL it's um one of the only Hadoop based systems that can run a traditional set of applications um um one very large Bank Bank located here in in San Francisco had code that they wrote years ago for uh a particular Financial Services application we're able to run it on our database in in their tests and almost all of the other databases that our modern query layers on top of Hadoop and Spark just don't have a sufficient amount of SQL yet we're really trying to look at what does it mean to be a modern rdbm M so we have to support a lot of capabilities and lastly um I think the biggest innovation in terms of computer science that spice machine has made is specifically in supporting distributed acid semantics being able to support transactions so what that means is that all of the acid properties atomicity and isolation and consistency and durability the all of these properties are maintained via a advanced mvcc architecture um the multiversion concurrency control architecture we use is standard snapshot isolation meaning every right to the database has a Tim stamp and no reader can see any rights that come after the timestamp of that reader so there's no locks on reading so the system just like an oracle system or a postgress system works it works on this distributed platform and we um leveraged some research that was originally done by Google in a system called percol percolator was a transactional layer on top of big table and that later was enhanced by uh Yahoo labs and in a system called Amid and we've brought some of the experts in from these projects and we're really excited to have extended that research in particular in handling really large transactions imagine somebody wanting to do um 2 million updates and have that be automically committed or rolled back we can handle that whereas some of the earlier research tried to keep all the Deltas in memory and sometimes that can be um challenging if the transaction is too large so that's some of the computer science behind our database and more practicalities though databases are used by people doing business intelligence business intelligence is um of course now a very robust industry there's tons of systems out there older systems like micro strategy um newer systems like Tau or Domo or something like that there are many many different systems even you know looking at Excel as one of these systems we are easily connected to any of these through odbc jdbc drivers and you know no matter what programming language you're you're developing in again odbc jdbc compliance gives you the opportunity to just write your application on on top of us we've got a great um group of advis ERS ranging from the um gentleman who ran the amp lab and the computer science department at Berkeley Mike Franklin that's where Spark was created Roger Bamford was the creator of Oracle rack one of the first employees at Oracle Maria namat was the founder and VP of engineering of times 10 one of the first inmemory databases and then later took over all inmemory databases at Oracle when Oracle purchased x 10 abanov Gupta was a co-founder of Rocket Fuel where I'm a chairman and he was one of the earliest deploy um deployments of uh he he engineered one of the earliest deployments of hbase at a global scale uh for adtech and being able to um perform incredible numbers of of AD bid impression evaluations in literally milliseconds and Ken Ruden is the um head of analytics at Google for search analytics and formerly had that role of head of analy ICS at Facebook and so great group of people lots of good press advanced management team um and again just summarizing and I'm going to open it up for questions what we're all about is building a traditional rdbms system it's a SQL system it's scale out you don't have to worry about what data flow engine is under the covers we worry about that in our Optimizer and essentially simplify the ability to make decisions in the Moment by you know getting rid of the need of B oltp from oap or having complex Lambda architectures so with that I'd like to open it up for questions so ma Appliance and an appliance who the to transform very good question so the question is is it an appliance and and if it is an Appliance how do we transform or or extract the data onto that Appliance and um first it's not an appliance this is all software there's no Hardware involved in specialized hardware for us um we deploy in a couple of different ways um first um you can download it and put it on your own clusters right um and it supports any of the major Hadoop distributions whether that's Cloud era mapar or Horton Works um and in some of them it's very easily deployed as just a little one push of a button called a parcel um that just puts uh spice machine jars out on every um box and it then you just connect you start it and connect to it um with regard to analytical processing um if somebody has a boatload of data that they would like to essentially bring in and analyze they do it in a variety of different ways with splice machine sometimes it's as simple as just pointing us to a Hadoop file that might have billions of of Records in it and we just ingest it with a bulk import mechanism that uses every Noe in the cluster to simultaneously import and we make sure indexes are kept consistent through our transactional integrity sometimes people will stream data in as it's being created through Kafka for example or another streaming mechanism my question in you have all theal so you have I see a question and say 40 so for that processes I understand do you still need to use it great question do you still need to kind of transform from a normalized representation maybe into a star schema or something like that and I say that that's definitely um use case specific because sometimes the parallelism that's afforded by our architecture is fast enough where you can keep the data in the operational framework and just query against that sometimes you do want to convert it over to a star schema to afford yourself some of the Simplicity that that gives you in query processing and and even just um helping the data scientists think that way you know it it's useful to do the transformation but the point is then it's it's really um it's you're getting rid of the extract right it's and you're getting rid of the load you just do the transformation so it's just not you know not El or ETL because it's all on the same platform so you do you do have different people doing different things so my question this is how this is from because across all data sources that we M databas file system as and [Music] have any yeah I can in a in a nutshell and there are many things different between the systems of how they work but the principal difference is that none of those systems can simultaneously do transactions none and if you have an application where you have multiple users hitting the applications at the same time with concurrency and you want to be able to make little reads and writes and changes and updates those systems will simply not work they don't have the mvcc that provides the transactional Integrity um many systems are available for analytics only but if you really do have just some a data set that maybe gets updated once a week you're going to perform some big machine learning analytic or some scoring function on it many different systems may be appropriate like drill or even you know columnar databases like vertica but if you have a living breathing database application where data is changing on you in a fair in near real time and you have concurrency going on you need an acid compliant database in addition to your analytical workload and that's um what we've attempted to focus on any other questions yes so spark there is great question so the spark um system is controlled by yarn and we can allocate you know um to the yarn processor um and we've also done something which is your your Insight is Keen there because um spark typically is architected to have one spark context that is aware of the resources across the whole cluster and we've engineered a unique um ability to use a remote uh spark contact so that any one of those nodes can um find uh the resources and dispatch the query to the spark context that's on one and not have to transfer the data to that particular node that has the spark context it's a modification to spark that we may even give back to the spark Community um if the spark Community believes that this is a useful um capability but it's it's Paramount for us to be able to use spark in that distributed multi-user way all right well thank you very much for your time this morning I appreciate [Applause] it