scala.bythebay.io: Mohammed Guller, Introduction to Big Data and Spark
Recording: scala.bythebay.io: Mohammed Guller, Introduction to Big Data and Spark
you good morning everyone so today I'm going to cover basically are two topics a short introduction on big data what exactly big data is and then we'll talk about one of the Big Data technologies spot since I want me to in 10 minutes I wouldn't be able to go in a lot of depth but i'll be around so if you have questions later on just grab me and i'll be happy to answer any questions you have so before we jump in into the presentation a quick summary about myself I'm working at a startup called class Pam and I have a dual role there and the principal architect for all the analytics and the machine learning components that we built and I'm also the engineering manager for one of the products that we are rebuilding prior to this I've worked for a bunch of startups including two of my own and my passion is basically building new products machine learning and big data analytics I also wrote a book on Big Data index with spark it was published earlier this year and it basically guides you through different out with different kind of analytics that you can do with spark it has a lot of examples it also includes a primer on scholar for people who are not familiar with Scala and it also covers some of the other Big Data technologies that typically people use with Scala of it spark a quick overview of class beam so we are a SAS product company we are not a consulting company and our product basically ingest multi structured data from any kind of machine or any device out there and allows our customers to do analytics on it all right so let's talk about Big Data now so as most of you know data has been growing at a very exponential rate over the last several years and so where is all this data coming from a small chunk of that data is basically the traditional enterprise application so publications like crm applications erp applications website activity tracking applications all right some of that data is basically coming through these traditional corporate applications that most businesses have the next big source is social media which includes basically apps like Facebook Twitter Linkedin Pinterest and PS I'm sure everybody in this room is using one of the probably more than one right so people are constantly posting updates posting videos posting pictures so that's what a big source of data that is basically we have been gathering over the last several years but the biggest growth is coming from this thing called Internet of Things almost 40% the data today is basically from IOT devices so what exactly is IOT it's basically network of devices that are embedded with software for collecting data in sending it over the internet so this is not just devices that allow you to access the internet but they're actually instrumented to collect data about how exactly that device is getting used what exactly are doing with that device so one example of that device is the smartphone that we all carry so this huge amount of data that Apple if you are carrying iphone for example Apple is capturing in terms of our exactly that device is getting used same thing with smart watches or anything that has the label smart typically is basically capturing lot of data in terms of how the device is getting used and typically the product manufacturers using that data to for a lot of different uses but the end goal is to make the user experience much better and according to some estimates that there will be five times more of these devices then humans by 2020 a subset of IOT is industrial IOT so these are basically devices that are a lot more complex that probably lot more complex functionality and generally lot more expensive but they also generate a lot more data I've shown some examples here so this includes modern manufacturing equipment automotive equipment medical equipment that you see in labs and hospitals data center devices and just to give us an example of typical flight write a single flight it will generate on 7 terabytes of data and there are like thousands of flights every day probably hundreds of thousands of flights every day so let's talk about okay so there's ton of data right now when people talk about Big Data what exactly is Big Data the three key characteristics are shown on the slide the 3 v's now people have added more ways over a period of time they're like four or five vs now but I want to focus on the 3 v's today first one is volume or by which I mean the scale of data so a lot of companies are collecting terrible of data every day actually and some of the big ones like the ones that have shown your google for example at 15 plus exabytes of data facebook has two plus exabytes of data and yahoo actually has 600 petabytes just on one Hadoop cluster and they have probably around 20 such clusters the second V is the velocity or the speed at which data is getting generated so it's not like traditionally how business applications are generating today the data is getting generated at a much much faster pace so the other name that a lot of times people will use is in addition to big data is fast data and I've shown some examples here so on what's up for example there are 21 million messages every minute on Facebook 5 million likes every minute the third way of big data is variety so traditionally most of the data that is getting generated was structured data that is stored in a relational database that's not the case anymore actually most of these data that's getting generated today is either unstructured or multi structured data so you can't really use traditional tools to process that so big data comes with big challenges and I've shown the two of the key challenges one is storage the other is processing and we'll talk a little bit more about each of these challenges so the storage challenges that traditionally organizations use San or nas devices stores data but unfortunately these are expensive devices so if you are storing like petabytes of data on these devices it will get a very very expensive the other issue is what kind of software you used to store that data right so they're basically that you store them as file so you store them in some kind of databases now traditionally people use relational database like my sequel postgres or Oracle or sequel server but those databases were not designed to handle this kind of data they can't handle the 3 v's of big data so that's the storage challenge the next challenge is processing how do you process of this data right when you're generating terabytes of data every day how do you process that one issue is the processing needs are a lot more diverse today than they used to be the right so traditionally most of the data was basically put into a database and people were analyzing them using traditional sequel statement so basically the usage was bi are not bi basically our traditional analytics but today most of the needs are beyond this is what you can do with sequel the other issues we want to be able to process that data in a reasonable amount of time now the definition of reasonable amount of time depends on the application for some applications it could be minutes for some applications it could be hours an option two examples here right so if example if you are building on the e-commerce fraud detection exact a shin you want to be able to detect whether the transaction is fraudulent or not within seconds right you can't take minutes or hours because if you that it's already too late on the other hand if you are building an application that's directing a detecting fraudulent insurance claim there you have a lot more time actually right so you could take hours that's fine so how much data i canna standard server process today it depends it depends on the definition of standard server it also depends what's the Rizla muhammad's your definition of reasonable amount of time so what standard today was probably are very high and server a few years ago and what's high-end today is going to be wrongly standard in figures right as with Moore's law hardware becomes cheaper and cheaper as time goes on the other thing is so let's say three years ago right ten years ago 100 GB would have been probably too much to handle in one one server but today most likely for most applications hundred gb is probably not a big deal you can handle it on one server right but as you start moving towards saw on your left it becomes harder and harder so what do you do right if you have terabytes of data or hundreds of terabytes of data that you want to process on a daily basis what you do or you have to option one is you can scale up the other is scale out so let's look into Modi or look in more detail each of this option in a scale of architecture basically what you do is you by powerful high end server so these are servers for example that may have hundreds of course or maybe 200 course in just one machine so maybe they have a lot more memory so that some high-end servers that can give you up to 10 terabyte of memory in just one box and they are basically architected for high performance but higher scale of architecture has its own disadvantages the first one is that generally these are hardware machines that have propriety architecture and as a result they are very expensive not everybody can afford them and the other issue is even if you can afford them they're still provide only a limited scalability right even if you have 10 terabyte of memory in a single box what what if you have hundreds of terabytes of data that you want to process still you can't do it with one box so the other option is scale out architecture in a scalar architecture what you do is you basically put together a cluster of computers to build one virtual supercomputer so you're pulling together the computer resources that's available on each node which is basically the CPU the memory in the storage and you build a powerful computer using that and still out architecture provides couple of benefit one is that it's relatively inexpensive so for the same processing Hart's power that you can get from a scale-up architecture you can get it from a scalar architecture at a much lower cost so that's one benefit the other benefit is that it allows you to scale up basically it's a conical to scale so you can start with the three node cluster which doesn't require a huge upfront investment right so in the scale-up architecture to buy one of those powerful boxes you would have to immediately make an investment of hundreds of thousands of dollars which not everybody can afford whereas with this architecture you could start with a small three node cluster which may not cost you more than five thousand dollars and as your need grows you basically keep adding more and more nodes and then you have more horsepower to process that data as well as stove that data however a scale-out architecture also has its own some challenges the first challenge is that how do you write applications that can take advantage of this distributed architecture so you need to be able to take advantage of the CPU that's distributed the stories that distributed and not everybody can do that it's a weight of problem it takes a lot of time to build this basically you need to be able to in your application chunk out the jobs right take your job built it into small pieces ship it out to all those different nodes and then coordinate whose then who's not done and then that then we bring this girl back he also need to be able to hide the failures because notes do crash hard disk crashes sometimes there might be neat network issue so you need to be able to handle all those things the other issues that typically when you write program that's running a single server the probability of that node having some hardware issues very very low typically less than 1% but when you hundreds of these nodes in a cluster the probability of any one node filling is pretty high which means you do have to tackle that in your application all right so let's talk about one of these solution technologies that was built to basically that allows you to use scalar architecture as well as I have been at the same time it addresses some of the challenges that we just discussed so what exactly is spark it's basically a fast easy to use general-purpose framework for processing large days large data sets using a cluster of servers so in other words is basically allowing you to implement scalar architecture but at the same time it's making it really easy to use and it provides are very easy to use of actually a very fast framework that you can use for a variety of different or data processing needs and we are going to more depth into each of those areas so the first thing it gives you is it abstracts distributed computing now as a programmer as a developer you don't have to worry about okay how do I ship my job to different nodes spark is taking care of that for you actually hides all the messy details of distributed computing for you you have to just focus on the business logic that you have to implement the other good thing is you write code and the same code exact same code works either on a single laptop and you can then once you are satisfied that okay it works the Lord you got the logic right you can then deploy it on a cluster which has hundreds of nodes and it for exactly the same you don't have to make any code changes the other good thing is it as it's highly scalable so you could have a cluster of let's say three nodes in the same code you can then deploy it on thousand nodes and automatically your code will scale because of all the things that spark provides to you its fault all range so even if let's say half thousand node cluster few notes crash you don't have to worry about that in your applications spark is handling it for you the next thing is fast it's extremely fast and they didn't have coded them in a different color is just to differentiate spot against hadoop mapreduce so prior to sparkle a lot of people were using Hadoop MapReduce now MapReduce also basically provides these first three things as abstracts tributed computing it's scalable its fault on in where spark really shines us these three things it's fast much much faster than MapReduce it's easy to use and it's flexible so let's talk about this speed angle first right so if your data can fit in memory spark can be up to 100 times faster than hadoop mapreduce but even if your data does not fit in memory it can be still 10 times faster and I've shown a bench bench mark here that was in long time ago actually basically Baron or machine learning algorithm called logistic regression on both using Hadoop MapReduce and spark and as you can see it took only less than a second done spark weather than Hadoop it took almost hundred and ten seconds so why exactly is it fast a lot of people ask me that two reasons one is that it allows applications to catch data in memory and that gives you gives you some speed advantage it but the other benefit is the other reason it's fast as this job execution engine it has a much much better job execution engine than what how does not reduce provides to you let's talk about the second thing that differentiates part against all out of map it's much much easier to use so basically as an application developer use park as a library and it supports four different languages it supports scala its supposed java python and our and then the core API which was the RDD API had almost like 80 + operators compared to the TV operators that MapReduce provides now this is kind of becoming more like the low-level API and with the newer versions probably you may not have the need to use that you'll probably use this data set of the data frame API which is a much higher level abstraction makes it even easier compared to the rd api which was itself lot lot easier than our MapReduce it also provides a sequel interface so if you don't try like Scala Java Python at great no problem you can just do all the data analysis using just sequel it comes with a couple of interactive development tool which again makes it easy to try out different things in an interactive environment so head comes with a spark shell that's basically built on the Scala shell and they're also a couple of notebooks that you can use if you are coming from the Python background and you're used to using notebooks the next thing I want to talk about the flexibility right so spark is basically not just meant for batch analytics of data you can use it for variety of different data and it you can use it for it comes with couple of libraries so the sparks equal spark streaming Emily band graphics these are the four main lab bridge that come with spark and you can use this to basically do not only batch analytics but interactive analytics for machine learning frog for graph analytics and four interactive analytics as well so one thing to keep in mind is the spark is basically a pure compute engine it does not have its own built-in storage so what sometimes people confuse it that okay can i replace my relational database with spark not really actually because typically a relational database as both a compute engine or the sequel execution engine as well as a storage engine spark does not have a built-in storage engine so this one benefit of that is that actually don't need to import data from other systems so if you already have data in some other system you don't have to import it into Spanish chicken process from spark for data from wherever it is directly the other benefit is that you can actually scale your compute nodes independent of your storage nodes right so you don't have to have take into account okay how much data I have and how much disk I need and based on that calculate how many sparknotes you want you can just have the sparknotes based on how much compute processing how much compute power you need alright so the good thing is part can process data from a variety of sources and I've grouped them into three categories the first one is basically distributed file system so it supports HDFS as well as Ennis on s3 and on these file system it supports a variety of file formats so your data could be in a pro format or power K or CSV or chase on the next box is basically the no sequel data store so if you are using one of this modern no sequel t testers like Cassandra or HBase or elastic search again you can use analyze that data using spike there's no need to import that data from that data store in to spark the third category is the traditional relational databases again any database that provides JDBC connectivity you can cross those data from the database using spot alright so this is the last slide so what is spark really ideal for right what kind of applications really benefit from spark so there are three classes are three classes of application the first one is application that have a really complex data processing pipeline so typically in the Hadoop MapReduce wall right if you have a complex processing pipeline you have to write multiple jobs and you have to synchronize those jobs and the data gets written into disc after every job so there's like right and then read write and read and you have to synchronize all that with Spy you don't have to do that right pretty much one job can have any number of stages and so you're not you're not being that penalty of having to write to disk at the end of every chop so it really spark is very very good for those kind of use cases the other use useful class of applications for spark is other applications that use my creative algorithms and options two examples here so they're like machine running for example you know typical machine learning application you make sometimes 10 passes over the same day Dorothy 100 passes over the same data so if you have to do that right where you're kind of right rating over the same data again and again that's where spark is really really good actually you'll see a lot of performance benefit for those kind of complications the third is a dork analytics so again since Park is fast right given the fact that it's park and fast and it gives you interacted all open environment if you want to do interactive analysis Park is very useful for that actually can analyze large amount of data so typically what will happen is let's say if you analyzing terabyte of data right the first time you issue your query it will be slow because it has to read that data from disk but then you can cash that data in memory and once it is cached all your subsequent queries will run much much much faster all right so that's a quick overview of big data in spark I don't know if you have any time left but I like I can take my question ok nobody and I think the spark arts park two point oh I'm sorry you got anything with spark are in version 2 point 0 what greater what do you mean by god I'm familiar with it but I've i personally use Parker with two dot oh is that the question now I haven't used and more on the scholar site so pretty much everything I do is with Scala thanks thanks [Applause] you