Devreal

Apache Spark, Apache Flink, and Apache Ignite: Where Fast Data Meets the IoT

Event: Apache Spark, Apache Flink, and Apache Ignite: Where Fast Data Meets the IoT

SF Spark: Denis Magda, Apache Spark, Apache Flink, and Apache Ignite: Where Fast Data Meets the IoT

Recording: SF Spark: Denis Magda, Apache Spark, Apache Flink, and Apache Ignite: Where Fast Data Meets the IoT

[Music] so thanks to my experience in big in distributed systems and in the IOT market I can share with you some of the use cases or how you can apply distributed frameworks such as spark flame can ignite to solve the cases for IOT applications and services how you can build an infrastructure that will be able to process all the workloads coming from your smart meters from your smart sensors from your edge device is deployed in factories or in plants and the agenda for today is pretty simple so first before we start diving into the into the kind of into the way how a problem can be solved let's take a look what are the demands that are coming from IO et empowered applications and services after that when we know these demands I would introduce a high-level software stack that you as an architect should keep in mind it's not the detailed one it's not will be a breakdown but it's something that some the main building blocks you will have in your IOT solution or project and after that we are going to drill down on every module of this stack we are going to look into device operating system or real time operating system that that is running on your edge devices and to make things clear where the edge device is any device such as smart meter sensor smartwatches everything else that is we can strain device and gathers and measurements information around your environment or around the field of the device applicability also as we will see there are different as additional layers that are used in this architecture such as data collection and enrichment layer and for sure we are going to use some data storage some distributed storage which I call hybrid transactional and analytical platform and finally when we come up is our HT platform we want to be sure that that platform gives us enough api's needed to build powerful applications on top of it and usually I end up every my conversation every my presentation with a demo and that meeting is not an exception for us so at the end of the of this talk I will show you how you can quickly create a simple application using a batch ignite apache spark and process transmit some data stream something to into the cluster and process it in real-time well as for the demands the first demand that is pretty crucial for the IOT space is real-time processing well yes the first one is real-time processing and here is really mean that for some of the applications in the Terry tier as real it's absolutely critical to be sure that they told the data all the measurements that delivered to your software stack can be processed and received in real time just imagine that your software is running on some auto named automatic driving car and that car is usually quit by variety different sensors and it's totally it's it's critical requirement to be sure that every measurement every information that is coming delivered by your sensors is processed in real time if you cannot meet this requirement you can get into a car accident and for sure we want to avoid this the next requirement is is how your platform is brought in terms of api's iut developers if we talk if you are talking about the back-end developers or high end developers the guys who are writing back-end software for your IOT solution they also want to use common API it's like SQL like key value operations and what's else specific for IOT is if you are if you if you invented some project that transmits your location your proximity and also you want to store the data in so-called geospatial format and it will be not even it will be beneficial if you play it from support geospatial data and API is out of the box the next thing is analytics so for some of the products it's it might not be requirement but for the rest when you want to understand what's happening to your system or what's happened to your smart meters deployed in some factory there you go a week ago you want to execute some analytical workloads and most likely you will be executing nowadays you will be executing them in parallel to the operation of workloads in parallel to the new information that is arriving to your software stack and considering this nowadays high viability becomes and became not just a nice feature to have but a sort of a requirement a lot of their projects they want to be sure that whatever happens to my system whatever happens to my architecture I want to be up and running unless you know the whole data center goes down but even if the whole data center goes down there is a solution right you can have several data centers that work in parallel and that's all about here'll ability and finally rivalry close to the cable high availability requirement we have simple scalability requirement that implies that first they want to invest in such a solution that allows us to scale out great Julian so probably at some point in time we are going to start with hundreds of meters smart meters deployed somewhere but all the time we want to scale out you want to deploy thousands and thousand one thousand meters and our software architecture has to react accordingly it it should be simple it has to be simple to scale out your existing software and basing on these demands actually we can take a look at this from these IOT software stack perspective starting from the bottom first if you're going to create your new IT solution like first you need to deploy some edge devices like smart meters somewhere and those smart meters they also have to be programmed the head you have to develop applications from them software and usually there are a special type of operating systems and KPIs that are developed for this sort of devices if your device is powerful enough and you can run a Linux operating system embedded Linux on it then you are good to go you can install Linux and you can use POSIX API is provided by the operating system but usually all these edge devices they are running on constrained Hardware and under the constraint hardware I mean the chipsets that have around would say let have those and kilobytes of RAM and a couple of megabytes of flash this is all the hair but that space that amount of memory is is enough you know to install in a real time operating system on it and to develop applications that will gather different metrics surrounding the smart meters and transmit all this information back to your distributed storage in bed to your software stack on top of this the next layer which is called data collection enrichment usually this layer interacts with your smart devices in this layer you're going to connect to a smart device devices pull information from there or the smart devices will put the information to this layer you can enrich this data meaning that you can take the raw data and transform it adding something useful adding some payload that is needed for your application and right after that right after you update it right after you receive the data you want to store it somewhere and for big deployments of other deployments that are supposed to scale out over the time we want to store this data in a distributed platform in the distributed database and we don't want in this distributed storage or database to be only you know SQL based I want to interact with this storage using another api's and that I needed for my application and this is why we are looking here into the so called H step play drums or new SQL solutions and that platform will be capable of storing the data in a distributed fashion and will provide a variety of api's that you add that your application can benefit from and today i am representing applied to ignite software foundation and we are going to discuss the software that perfectly well feeds to every layer depicted on this diagram and all these software belongs to a budget the bottommost layer here is we will be I will briefly introduce you to Apache minut it's a real-time operating system that is this is actually the first and the only one at the moment a real-time operating system that is available under Apache to that as for the next layer data collection and enriched here is we have much more options for the date you know for the initial data processing and receiving in transformation we can use part we can use Kefka we can use link and today we are going to briefly look at CAF we are going to sorry to look at spark and Flint and I will show you during the demo how you can wire up together ignite and spark and ignite is going to be used as our distributed computational platform and our distributed storage on top of once you use ignite for this purpose your application will be empowered by variety of different TV eyes and today I am going to show you some of the api's that are useful for the iot use case we're not going to explore the whole platform because it would take a lot of time I just want you to learn something that is specific for the IOT applications okay now let's take a look at air we building block from here in more details so first I applied to my new what's so special about this operating system first this is open source real-time operating system that really can run on a constrained devices on the devices or microcontrollers that are empowered by cortex-m processors CPUs for instance and if you to give some you know basic understanding on this a microcontroller usually it's small chip that has a CPU Ram flash and sometimes networking components embedded and you usually take this chip in this small chip and I build your smart watches your smart meters using that hardware and once you construct it once you assemble to your hardware you can flash Apache module there and use the rest of api's its supports for instance you can write standard and networking applications you can communicate to your back-end system using Wi-Fi bluetooth and networking protocol em utilize tcp/ip or OTP at the same time when you need to update your software first you apply to minutest many are the real-time integration systems they guarantee that the data that is going to be flashed the software that is going to be flashed on into the hardware that will be secured so you can define some security policies so double check that nobody else rather than you can update the software on your hardware and speaking about the updates more than three alternate rating systems they do not you know force you to come to your factory or to the plant where your smart meters that deployed or Apple doesn't require you you know to deliver all its smartwatches back to the office if they want to update the hardware software all this stuff can be done remotely and you ISM as an IT architect or ret developer he also wants to install Apache minut you can update Apache my new tree ocean in the ocean of your software remotely from your office the next layer data collection and enrichment layer that's layers so as I said it interacts directly interacts with your edge devices so usually you can connect to your edge devices using different networking protocols like Wi-Fi Bluetooth depending on the scenario and once you connect to your devices you can utilize spark streaming functionality or flowing streaming functionality to receive initial data from there you can as you know you can exit you can scale out your spark question for instance depending on the workload coming from Europe smart meters and later what once you receive this data you can process this data and in parallel enrich this data with additional information and right after that when the initial pre-processing processing of the data is done you can move it on to your final storage to your distributed system where the data will be spread out across a cluster of machines and on top of that system you can write additional applications that can do some useful thing with that okay and if to talk about that platform about ignite as our each tape platform as our distributed storage and computational database so actually a batch ignite is in memory computing platform that is durable always available and provides powerful sequel key value and are the products in the api's so at the core of this system lies so-called memory centric storage the memory centric storage implies that the RAM main memory is treated not education where its treated as one primary storage at the storage where your data as well as indexes will be kept by ignite at the same time the memory centric means that you can still use disk as your secondary memory tier just for the sake of the durability if you want to be sure that if some bad happens to a cluster imagine that your whole data center goes down and you just you for sure don't want to move the data but if you enable this persistence if you enable ignite native persistence for sure there you can be sure that after the cluster to start you will not miss a piece of your data all the data will be persisted on disk and you can start working with your data right away so ignite native persistence is tightly integrated with the overall memory architecture and usually it stores the super set of data on disk and you can have as much data in RAM as you can afford so and at the same time you can perceive the data on third party data bases so storages like NoSQL databases like MongoDB Cassandra or databases it's up to you but there are some of the differences between ignite Native persistence and third party persistence so if you if you if you have if you are interested in this then I'll be happy to elaborate more on this at the end of our talk but okay and all this storage is distributed right so all this memory is available across the cluster of machines you have and once you have your all the data and indexes in RAM you can leverage from the variety of API is provided by ignite like SQL key where you might reduce framework which is which we call compute Grid machine learning and other capabilities if - if to look into this distributed storage in reality this distributed storage is a key value database so the same hash table so we all know how to work with hash tables right so every hash table consists of several buckets and here is we are dealing with the same hash table with the only difference that every bucket is belongs to a specific cluster node specific last year machine or ignite and we call in ignite we call these buckets as partitions some other products call these kind of buckets as shrubs when people use chard in technology to spread out the data and for sure once you spread out this data evenly across the cluster of machines you automatically get you can automatically lower balance their workloads so every API call that comes to your application depending on the data locality that API call can grow to one twice your note or to another and this is how you can scale out your overall data storage and this is how you can so called scale out your workloads by just adding you nodes to the cluster whenever it's needed the next so speaking about the api's so the since that's the key value storage probably the fastest way possible to interact with the key value buy storage is to use key value operations so if you know an ID of your let's say an ID of your smart meter and you just want to get information about this smart meter for instance where it's deployed in which geographical location then you just can take this idea of the smart meter and set the query to your question and then we can take this ID and we can identify precisely where this information is located because for every ID we can calculate a primary node that stores the data for this type of record and if you execute millions or hundreds of thousands of such queries in real time you can see that all these queries will be sent to different nodes and this is how you can distribute your workloads by using simple key value operations the other API that is useful for IOT applications is SQL it's required when unive on to execute some advanced queries like you want to see what was the highest temperature reported by your smart meters or where what was how many smart meters failed throughout the previous week all this information is easy easy to retrieve using SQL language right so and if you don't want to invent your own language ignite offers you full-fledged sequel PI's on top of it and even more so it's quite often ignite is used as a just full-fledged sequel database simply because you can execute not only standard select queries but you also can update data using operations like inserts abate deletes if you want to configure your cluster data sets you can use create table commands you can define indexes and redefine index history in real time in all this operations can be done not only from Java but near to C++ libraries that are provided by ignite out of the box but you connect to your to last year from your favorite tool from your favorite language using ODBC or JDBC drivers this is what we provide here the next interesting thing is computed so actually compute gritties ignites a MapReduce framework so the same framework you accustomed to use in spark Oh in Hadoop so the idea is the same you write your computational logic and you issue you send this computational logic to your question and this end depending on the type of this computational logic it can be just randomly distributed across the cluster of machines as we see here and if you apply different load balancing techniques for this logic for different calculations then ignite can you know prefer and the loaded cluster nodes for new jobs calculations rather than overloaded nodes in your question speaking about spark so if you use uniting spark together for your IT solutions as spark usage is not limited to this dreaming or apart it's not limited only to the data processing and enrichment step if you like if your kind of advanced spark user and you prefer to keep working with spark applies ignite is integrated with spark so you actually what you can do you can connect from your spark applications to your underline ignite cluster and work with the data stored in ignite using spark RDD CPI's for instance so I will show you today how we can connect to the to ignite cluster from spark get access to our data sets stored in the ignite cluster from spark share targets and how we can you know work with this data using that already known spark API yeah say what is it ignite database ignite so ignite can be deployed on previous or in any cloud environment it doesn't matter so you can go you can go in deploy ignite on AWS Microsoft Azure or Google compute engine well you can as I show you today you can deploy your ignite question on you look a laptop it's up to you now you're just there are two things first there are ignite is available on Amazon Marketplace so you can installed from there or you can just go ahead and download ignite from ignite site and upload it to your Amazon Cloud and set it up there is no so the question is designated supports as three storage so there is no eNOS support out of the box but if you want to use s3 as your persistence layer there is a straightforward you know like five methods API that you need to implement that's it if you want to use f3 as your storage but most likely you don't need to use f3 you can use for instance if you use a patch ignite persistence and you need to do some backups backups copy of your data then s3 is a good like is a good place to store your backups something like that but there are a lot of options options how you can use ignited the poet on AWS good disappear and finally for this scenarios one your IT solution when your IT applications need to do some more than just advanced lookups using SQL or it's even not enough for you to utilize our MapReduce framework then you can take a look at the machine learning grid so it's a set of machine learning API so that are being developed on top of ignite and actually there is a reason why we decided to initiate this development because we talked a lot to our customers and users and they shared the idea you know guys we already use your computer with your MapReduce framework and we use your distributed storage capabilities and built a built our own machine learning applications on top of it why don't you go ahead and you know implement your own why don't you simplify that why you don't use facilitate this part from your side and develop this API is on a patch ignite side and now if we have some machine learning into the assets in our party net community that already developed a distributed core algebra framework on top of a patch ignite distributed storage and they are working on a variety of different essential machine learning algorithms such as k-means clustering linear and logistic regressions decision trees and so on so forth so this is something that is new that component is available in beta and here is I'm just mentioning it because if there is someone of you because machine learning is still required in Houston IOT and if there is someone of you who is passionate about machine learning if someone who wants to implement some distributed version or well-known algorithm just stop by Apache net community we will welcome you and we will help you to make you successful and you will be able to contribute to this library making something you useful and building your skills in this direction well that was I think that that isn't enough about api's about capabilities so these capabilities I just picked on if those capabilities better related to IOT for IOT scenarios and speaking about and and if we it's obvious why we should use spark of link right for this sort of scenarios but why should we rely on apache ignite zone speaking about apache ignite so it has a lot of uses use cases most of them fall into the financial services market because of the guarantee of strong consistency and transaction guarantees that are provided by this platform but at the same time there are in mind is so some of the company's sites such as Silver Spring networks clean works that already use Apache ignite exactly for the IOT scenarios using the patch ignite as a distributed storage and computational platform for instance if to quickly look at this use case that was shared by Sylvan Sprint network with a patch ignite community what they did they are using grid gain cluster grid gain is an enterprise so grid gain I use this apache ignite builds on top of a patch ignite and provides enterprise level features on top of the open source functionality that is available for everyone if you download the patch ignite so if you if to think of grid gain it's the same apache ignite with extra features on top and still the sprinkle networks uses grid gain in production because of the additional enterprise features required and what they did there is a special platform developed in-house and that platform helps to track information about many electricity pillars deployed nationwide and depending on the world different regions or states of the United States they can move the overall memory the the overall electricity usage from one region to another if they face the situation when in one region we see overconsumption of the electricity but in the other region we see that the electricity is is not consumed that well they can move all this stuff over the wires and they have a kind of an easy problem at some point of time they could not meet the SOS and the architecture was built on on a classical relational database that was used as a storage and served and processed of the api's of the course but once they are moved to the grid gained in the patch ignited the distributed cluster they were able to deploy million thousand of millions meters and all these electricity pillars nationwide and process all the data and all the measurements in real time this is what they could achieve with this type of architecture and now the part that I like more let me show you quickly a quick demo so I want to demonstrate how we can how you can quickly start with a batch igniting spark and I'll ask Laura to facilitate me with this because I need both hands so your theme here is a cave I prepared a simple application I'm going to start a couple of ignite cluster nodes in my local laptop it's really straightforward operation all I need to do is to call this API method and pass a batch ignite configuration this configuration includes just basic parameters such as IP addresses of my cluster nodes so that they can find each other and form a single class through machines okay let me start several nodes well we have the first notable running we can see that presently there is only one node plus your note in our in my local cluster and here is we are waiting for the second one to join this party okay and now we have to cluster notes I kept to Apache ignite notes on my local laptop running and the next thing would be the fall let me just first you how you can monitor the state of your question for that purpose Apache ignite community developed a special management in monitoring - it's called the PI chick net web console this one is named grid gain because you can take this apply check net web console and deploy on your own hardware and you can change you know the loggers all the parts you need so this console is deployed on grid gain site you can also use it for your testing purposes or even in production if you're ready to do this I'm going to use it use it for the sake of the demo so here is let me double check that my I can connect to my local cluster from this console yes so I see my talk last year notes that are running on my laptop there as I'm not going to use the persistence of the data will be in memory I don't care and also once I started my cluster notes I created two caches so in terms of ignite we use the term cache it's not just the cache of data in memory it's just it's the same as a table in relational databases world but in the distributed world in the world of Apache ignite a cache is your unique table distributed table where you store your data and we are going to have two caches first one will store information about different sensors deployed somewhere and the temperature reported by these sensors okay now it's time you know to start doing something useful with that and let me show you how you can walk work with your cluster from spark and the spark application developer so here is my simple spark application what I am doing when studying I'm preparing spark configuration and creating streaming context I'm going to use spark streaming framework and after that I need to do some basic configuration parameters needed for spark to talk to ignite and here is you can see to those two caches that are already configured in our cluster and I am going to connect those caches by means of spark shared by spark RDD api's this is what I am doing in these two lines and finally the next we are going to do two things the first thing is we are going to promote our sensors cache with some damage data for instance what I'm going to do I am going to create several sensors and you know install deploy them in various geographical locations randomly picking a latitude and longitude once i've prefilled filled in this information i'm going to save this data into the apache ignite by calling this save pairs method that's the standard method available for you on our DD api's once this is done I'm going to receive some sample temperature measurements from this localhost port number I'm going to process this measurements that will be measurements of the temperature reported by one of the sensors created before and once I do all these transformations such as I want to add a timestamp to every measurement I will I'm going to say if all this data transmitted all this data constantly constantly to my Apache like last year in real time so spark is going to receive the data process it transform and send it right away to a patch ignite question and once and while this data is going to be processed and sent to a patching night in real time I want to execute some I just want to see what's going I want to execute for instance select queries over the data that has been streamed into my apache ignite cluster and that query we're still using a pipe this part we are going to for instance series I want to see what was the maximum temperature how many sensors we have in my class that reported temperature in this range from seventy to hundred degrees that's a simple query okay let's start this simple application so now spark is doing his basic set up routines and right after that the spark is going to connect to a batch ignite using spray a special client connection this is why you see ignite logo here so spark is connecting to ignite using ignite special ignite this bi but you just need to do just valid configuration you don't need to do anything else just to provide configuration where your Apache light cluster is that's so what you need to do okay so you connected him here as we see some strangers this search errors report says us let us say us that we cannot we do not receive anything from this localhost port number no any temperature measurements are transferred to us and but to do this there is another simple Java Danny application so instead of you know bringing any edge devices here I wanted to simplify the test cross so I just created a simple java application that opens up this server circuit connection and once our spark connects to this application that which in real world will be a mesh of your sensors deployed somewhere we are going to create different temperature measurements random measurements for random sensors deployed in our cluster and write them into this circuit connection and then this part will do this will you'll do the rest of their work okay let's start this application okay we see that the spark was connected to the cluster and now if you scroll down you should see that that exception disappeared and this is the result set reported by our select query so the data is being - received and processed in real time processed by spark sent to my local epigenetics last year and this query is executed over the data that is an apache ignite this is how many sensors the IDS of the sensors that reported the temperature in the range of 70 and 114 gains at the same time you can do some more advanced lookups here as we can see that our cluster data is being is being grown so we keep putting more temperature data into my cache what I can do it now I just can do some more advanced queries over the data I have in my class sir so for instance let me see what would the maximum temperature reported so far by my sample application and it turns out to be that the maximum temperature was 1 109 for you guys or in Apache ignite you can execute not just simple you know sequel queries you can also execute distributed joins joining the data that is also stored in different cluster machines and using this type of join I want to find out the maximum temperature reported by the sensors located in the boundaries of the United States so these are the boundaries of the continent continental part of the United States and we are going to take a look at what's going on with the temperature in this part of the world ok we also can keep executing this query for every two seconds and for now we can see that the maximum temperature reported so far was 97 fring grades okay actually that's all for the GMO and before we before we move on to the Q&A session let me advertise you there in memory computing sign it it's a conference that is that will be that that will take place in San Francisco in the end of October and here is if you if you actually want to learn more about different in memory technologists about distributed systems not only about the pledge ignite for sure there will be many plenty of folks who are developing a batch ignite but also you can learn a lot about different solutions of vendors such as hazel castor ADA's mammoths Cuervo GB so if you want to discover more about all these distributed systems and memory technologies I highly encourage you to stop by this conference and actually you can check out whether you are like you're not a grid gain is roughly in two tickets per week for this conference so I'm going to share this presentation and then after the talk you can go to this link and try to sign up for this raffle and finally thanks to the Lord who helped me out today so if you want today we covered you know all the features of Apache ignite that are available in open source version but if you want to learn more how to gain how to set up datacenter replication or how to set up advanced security how to do data backups using a patch ignite question then you can discuss all these questions and topics Monsieur so having said that thanks guys for attending and now I'm up in to your questions [Applause] thank you in the temple [Music] you created the RTT or one of those to say let's see what happens yes [Music] it's a good question so once you call this method in our DD safe we are going to we are going to use special type of ignite on streaming technology that will take all the tuples started this oddity and we'll push them into the question once the data is stored in the United classic it will be garbage collected by Java Virtual Machine via spark application is running and also when you want to get this data back oh for instance when I when I I showed you how you can execute sequel queries from spark that sequel query was executed over the data that is in ignite so spark did not go to ignite it did not take the data from ignite it didn't reload the data back to spark so all the queries operate over the data that is already in ignite so you don't need to move the data back and forth from ignite to spark so yeah so ignite here is is sparked reads ignite here as a distributed storage such as Hadoop or something like if the only difference here is that for instance Hadoop is a disk based storage while ignite is used as memory centric storage that can keep data boss in memory in ramen on disk [Music] no no we do not duplicate the data so once the data is pushed to ignite it will not be stored in spark the yeah the data is always sitting now in this condition any other questions oh yeah if you crave so I'll be here around so if you want to discuss something you can talk to me I can give you more details on all the components we briefly discovered so far or you can share more details on use cases related to the IOT or different other markets or you also can talk to my colleague to learn more about grid gain use cases thank you [Applause] [Music]