1"use strict";(self.webpackChunkdsa_blogs_site=self.webpackChunkdsa_blogs_site||[]).push([[5911],{3905:(e,a,t)=>{t.d(a,{Zo:()=>d,kt:()=>u});var n=t(7294);function r(e,a,t){return a in e?Object.defineProperty(e,a,{value:t,enumerable:!0,configurable:!0,writable:!0}):e[a]=t,e}function o(e,a){var t=Object.keys(e);if(Object.getOwnPropertySymbols){var n=Object.getOwnPropertySymbols(e);a&&(n=n.filter((function(a){return Object.getOwnPropertyDescriptor(e,a).enumerable}))),t.push.apply(t,n)}return t}function s(e){for(var a=1;a<arguments.length;a++){var t=null!=arguments[a]?arguments[a]:{};a%2?o(Object(t),!0).forEach((function(a){r(e,a,t[a])})):Object.getOwnPropertyDescriptors?Object.defineProperties(e,Object.getOwnPropertyDescriptors(t)):o(Object(t)).forEach((function(a){Object.defineProperty(e,a,Object.getOwnPropertyDescriptor(t,a))}))}return e}function i(e,a){if(null==e)return{};var t,n,r=function(e,a){if(null==e)return{};var t,n,r={},o=Object.keys(e);for(n=0;n<o.length;n++)t=o[n],a.indexOf(t)>=0||(r[t]=e[t]);return r}(e,a);if(Object.getOwnPropertySymbols){var o=Object.getOwnPropertySymbols(e);for(n=0;n<o.length;n++)t=o[n],a.indexOf(t)>=0||Object.prototype.propertyIsEnumerable.call(e,t)&&(r[t]=e[t])}return r}var l=n.createContext({}),p=function(e){var a=n.useContext(l),t=a;return e&&(t="function"==typeof e?e(a):s(s({},a),e)),t},d=function(e){var a=p(e.components);return n.createElement(l.Provider,{value:a},e.children)},c={inlineCode:"code",wrapper:function(e){var a=e.children;return n.createElement(n.Fragment,{},a)}},h=n.forwardRef((function(e,a){var t=e.components,r=e.mdxType,o=e.originalType,l=e.parentName,d=i(e,["components","mdxType","originalType","parentName"]),h=p(t),u=r,m=h["".concat(l,".").concat(u)]||h[u]||c[u]||o;return t?n.createElement(m,s(s({ref:a},d),{},{components:t})):n.createElement(m,s({ref:a},d))}));function u(e,a){var t=arguments,r=a&&a.mdxType;if("string"==typeof e||r){var o=t.length,s=new Array(o);s[0]=h;var i={};for(var l in a)hasOwnProperty.call(a,l)&&(i[l]=a[l]);i.originalType=e,i.mdxType="string"==typeof e?e:r,s[1]=i;for(var p=2;p<o;p++)s[p]=t[p];return n.createElement.apply(null,s)}return n.createElement.apply(null,t)}h.displayName="MDXCreateElement"},3094:(e,a,t)=>{t.r(a),t.d(a,{assets:()=>l,contentTitle:()=>s,default:()=>c,frontMatter:()=>o,metadata:()=>i,toc:()=>p});var n=t(7462),r=(t(7294),t(3905));t(4996);const o={slug:"spark-explained",title:"Apache Spark Explained",authors:["parham"],tags:["apache spark","cloud dataproc","aws emr","explained","apache spark explained","data engineering","data science"]},s=void 0,i={permalink:"/blog/spark-explained",editUrl:"https://github.com/datastackacademy/dsa-blogs-site/tree/main/blog/2020-12-03-spark-explained (1).mdx",source:"@site/blog/2020-12-03-spark-explained (1).mdx",title:"Apache Spark Explained",description:"A quick deep-dive into Apache Spark, the most popular distributed data engineering tool.",date:"2020-12-03T00:00:00.000Z",formattedDate:"December 3, 2020",tags:[{label:"apache spark",permalink:"/blog/tags/apache-spark"},{label:"cloud dataproc",permalink:"/blog/tags/cloud-dataproc"},{label:"aws emr",permalink:"/blog/tags/aws-emr"},{label:"explained",permalink:"/blog/tags/explained"},{label:"apache spark explained",permalink:"/blog/tags/apache-spark-explained"},{label:"data engineering",permalink:"/blog/tags/data-engineering"},{label:"data science",permalink:"/blog/tags/data-science"}],readingTime:10.055,hasTruncateMarker:!0,authors:[{name:"Parham Parvizi",title:"Software Architect",url:"https://github.com/parvister",imageURL:"./../static/img/team/parham.png",key:"parham"}],frontMatter:{slug:"spark-explained",title:"Apache Spark Explained",authors:["parham"],tags:["apache spark","cloud dataproc","aws emr","explained","apache spark explained","data engineering","data science"]},prevItem:{title:"The Time is Now",permalink:"/blog/the-time-is-now"},nextItem:{title:"What is Data Engineering? How is it different from Data Science?",permalink:"/blog/what-is-data-engineering"}},l={authorsImageUrls:[t(4427).Z]},p=[{value:"Overview",id:"overview",level:2},{value:"The What/Why/How",id:"the-whatwhyhow",level:2},{value:"The What?",id:"the-what",level:3},{value:"The Why?",id:"the-why",level:3},{value:"The How?",id:"the-how",level:3},{value:"Getting starting with Pyspark",id:"getting-starting-with-pyspark",level:2},{value:"SparkContext: how to interact with data",id:"sparkcontext-how-to-interact-with-data",level:3},{value:"Resilient Distributed Datasets (RDD)",id:"resilient-distributed-datasets-rdd",level:3},{value:"Filter data",id:"filter-data",level:3},{value:"Group by key",id:"group-by-key",level:3},{value:"Conclusion",id:"conclusion",level:2}],d={toc:p};
1function c(e){let{components:a,...t}=e;return(0,r.kt)("wrapper",(0,n.Z)({},d,t,{components:a,mdxType:"MDXLayout"}),(0,r.kt)("p",null,"A quick deep-dive into ",(0,r.kt)("strong",{parentName:"p"},"Apache Spark"),", the most ",(0,r.kt)("strong",{parentName:"p"},"popular")," distributed data engineering tool. "),(0,r.kt)("p",null,"What is Spark? Why is it so popular? When and how to use it?"),(0,r.kt)("p",null,"Learn the difference between the sub components (RDDs, DataFrames, SQL, Streaming, ...), setup ",(0,r.kt)("strong",{parentName:"p"},"PySpark")," , and learn how to write\nSpark transformations using Python and Jupyter Notebook."),(0,r.kt)("h2",{id:"overview"},"Overview"),(0,r.kt)("p",null,(0,r.kt)("a",{parentName:"p",href:"https://spark.apache.org/"},"Apache Spark\u2122")," homepage says:"),(0,r.kt)("blockquote",null,(0,r.kt)("p",{parentName:"blockquote"},"Apache Spark\u2122 is a unified analytics engine for large-scale data processing.")),(0,r.kt)("p",null,"Apache Spark is unanimously the #1 ",(0,r.kt)("strong",{parentName:"p"},"Distributed Engine")," in the World. It's able to run parallel data processing and analytical\nworkloads on very large dataset. It is very versatile and provides APIs in a variety of languages such as Python, Scala, Java, and SQL.\nIt can run on a standalone machine, a cluster of machines,\n",(0,r.kt)("a",{parentName:"p",href:"http://docker.com"},"Docker")," or ",(0,r.kt)("a",{parentName:"p",href:"https://kubernetes.io/"},"Kubernetes")," containers, or on the\n",(0,r.kt)("a",{parentName:"p",href:"https://cloud.google.com/dataproc"},"Cloud"),"."),(0,r.kt)("h2",{id:"the-whatwhyhow"},"The What/Why/How"),(0,r.kt)("p",null,"Before we get into introducing Spark, let's take a look at the basics:"),(0,r.kt)("ol",null,(0,r.kt)("li",{parentName:"ol"},"What is Spark? And what are its components?"),(0,r.kt)("li",{parentName:"ol"},"Why is Spark the most prominent tool? And when should we use it?"),(0,r.kt)("li",{parentName:"ol"},"How is Spark used? How does it achieve parallel processing?")),(0,r.kt)("h3",{id:"the-what"},"The What?"),(0,r.kt)("p",null,"As mentioned above, Spark is the ",(0,r.kt)("em",{parentName:"p"},'"Unified analytics engine for large-scale data processing"'),". Spark is mainly a\ndata processing engine. It does not store data itself but it enables large scale processing of data stored on\ndistributed systems such as Hadoop and Cloud. It provides a unified API in Java, Scala, Python,\nSQL, R, and other libraries which abstracts away the complexity of Distributed Data Processing and makes it accessible to all developers."),(0,r.kt)("p",null,"Spark is used by both Data Scientist and Data Engineers due to its expansive APIs which covers Python, SQL, and\nML libs. Most developers find Spark very easy to use since it build on familiar concepts of SQL and DataFrames."),(0,r.kt)("p",null,"Spark is very feature-rich and consists of many sub-components. This could be confusing for people who are just beginning to\nlearn Spark. Here, we aim to quickly explain each component and their use:"),(0,r.kt)("ul",null,(0,r.kt)("li",{parentName:"ul"},(0,r.kt)("p",{parentName:"li"},(0,r.kt)("strong",{parentName:"p"},(0,r.kt)("a",{parentName:"strong",href:"https://spark.apache.org/docs/latest/rdd-programming-guide.html#rdd-programming-guide"},"Spark RDDs"))),(0,r.kt)("p",{parentName:"li"},"This is the legacy core component of Spark. It stands for Resilient Distributed Dataset (RDD). It breaks\nup datasets into smaller distributed resilient chunks and provides a series of\n",(0,r.kt)("a",{parentName:"p",href:"https://spark.apache.org/docs/latest/rdd-programming-guide.html#transformations"},"transformation methods")," which makes\ndata processing accessible and scalable. These transformation methods (API) enable developers to provide custom functions in\nPython (or Scala, or Java) to map, filter, join, merge, and store data."),(0,r.kt)("admonition",{parentName:"li",type:"info"},(0,r.kt)("p",{parentName:"admonition"},"RDDs are the lowest level of interaction with Spark. While they provide a great deal of functionality, they are typically overshadowed\nby newer and more feature-rich components such as Spark SQL and DataFrames."))),(0,r.kt)("li",{parentName:"ul"},(0,r.kt)("p",{parentName:"li"},(0,r.kt)("strong",{parentName:"p"},(0,r.kt)("a",{parentName:"strong",href:"https://spark.apache.org/docs/latest/sql-programming-guide.html"},"SQL & DataFrames"))),(0,r.kt)("p",{parentName:"li"},"Spark SQL provides mix use of SQL and familiar DataFrame APIs (from Pandas) to query and transform structured data.\nIt stores and loads data from a ",(0,r.kt)("a",{parentName:"p",href:"https://spark.apache.org/docs/latest/sql-data-sources.html"},"wide range"),"\nof data formats such as CSV, JSON, Parquet, ORC, Hive, and the Cloud, enabling\ndevelopers to merge data across all these sources. It further allows developers to write\ncomplex custom transformation logic using ",(0,r.kt)("a",{parentName:"p",href:"https://spark.apache.org/docs/latest/sql-data-sources.html"},"DataFrame APIs"),"\nalong their SQL code."),(0,r.kt)("p",{parentName:"li"},"Spark SQL and DataFrames are the ",(0,r.kt)("strong",{parentName:"p"},"most widely")," used component of Spark.")),(0,r.kt)("li",{parentName:"ul"},(0,r.kt)("p",{parentName:"li"},(0,r.kt)("strong",{parentName:"p"},(0,r.kt)("a",{parentName:"strong",href:"https://spark.apache.org/docs/latest/streaming-programming-guide.html"},"Spark Streaming"))),(0,r.kt)("p",{parentName:"li"},"Spark Streaming is an extension of the core Spark API that enables real-time data processing of high-throughput\nsources such as apps, twitter, social-media, websites, or stocks.\nData can be ingested from streaming sources like Kafka, Google Cloud Pub/Sub, AWS Kinesis, or simple TCP sockets.\nData is processed using transformation logic expressed in high-level functions like ",(0,r.kt)("inlineCode",{parentName:"p"},"map"),", ",(0,r.kt)("inlineCode",{parentName:"p"},"reduce"),", ",(0,r.kt)("inlineCode",{parentName:"p"},"join")," and ",(0,r.kt)("inlineCode",{parentName:"p"},"window"),".\nFinally, processed data is pushed out to file systems, databases, and live dashboards.")),(0,r.kt)("li",{parentName:"ul"},(0,r.kt)("p",{parentName:"li"},(0,r.kt)("strong",{parentName:"p"},(0,r.kt)("a",{parentName:"strong",href:"https://spark.apache.org/docs/latest/ml-guide.html"},"Spark ML or MLlib"))),(0,r.kt)("p",{parentName:"li"},"MLlib is the Machine Learning library of Spark. It provides a wide range of ML algorithms such as regressions, classifications, and\nclustering methods on top of the familiar DataFrame APIs.")),(0,r.kt)("li",{parentName:"ul"},(0,r.kt)("p",{parentName:"li"},(0,r.kt)("strong",{parentName:"p"},(0,r.kt)("a",{parentName:"strong",href:"https://spark.apache.org/docs/latest/graphx-programming-guide.html"},"GraphX"))),(0,r.kt)("p",{parentName:"li"},"GraphX is a new component in Spark for graphs and graph-parallel computation. It provides data processing and data analytics\non data that is stored as graphs. Graphs are commonly used in social networking platforms to store people's relationships to\none another. GraphX enables developers to traverse and explore this data at large scale."))),(0,r.kt)("admonition",{type:"info"},(0,r.kt)("p",{parentName:"admonition"},"Most prominent components used by Data Engineers are ",(0,r.kt)("strong",{parentName:"p"},"Spark SQL"),", ",(0,r.kt)("strong",{parentName:"p"},"DataFrame APIs"),", and ",(0,r.kt)("strong",{parentName:"p"},"Streaming"),".")),(0,r.kt)("h3",{id:"the-why"},"The Why?"),(0,r.kt)("p",null,"Spark has been the most dominant tool due to a few simple facts:"),(0,r.kt)("p",null,(0,r.kt)("strong",{parentName:"p"},"Spark is Fast:")),(0,r.kt)("p",null,"Spark is much faster than some its predecessors such as MapReduce. Spark does most of its computation in memory while providing\nscalable, reliable, and fault-tolerant data processing across machines. It also has much less overhead to launch and coordinate jobs\non a cluster, making it more suitable for applications that require low latency to run."),(0,r.kt)("p",null,(0,r.kt)("strong",{parentName:"p"},"Easy to use:")),(0,r.kt)("p",null,"Spark provides APIs in a wide range of familiar interfaces such as Python, Scala, Java, SQL, R, and DataFrames. This makes it very easy\nto adapt for developers from various backgrounds. It's also very easy to mix these workloads together."),(0,r.kt)("p",null,(0,r.kt)("strong",{parentName:"p"},"Runs Everywhere:")),(0,r.kt)("p",null,"Spark can run on a single standalone machine, a cluster of machines (Mesos), Docker or Kubernetes containers, or on the Cloud. The same\ncode that is developed on a single machine can run on a cluster of machines, providing scalability. All Cloud\nvendors provide a Spark Service on their platform."),(0,r.kt)("admonition",{type:"info"},(0,r.kt)("p",{parentName:"admonition"},"It's also fair to point out why Spark is ",(0,r.kt)("strong",{parentName:"p"},"NOT")," always used on the Cloud. Spark services are typically expensive to run on the Cloud.\nOther ",(0,r.kt)("em",{parentName:"p"},"serverless")," services such as Cloud Functions and Cloud Containers are often preferred due to their cheaper cost. But if you need\ndata processing at a large scale, there's arguably no replacement for Spark and people pay the price.")),(0,r.kt)("h3",{id:"the-how"},"The How?"),(0,r.kt)("p",null,"Spark loosely works based on Scatter/Gather distributed processing techniques. Large data sources are divided up into smaller, more\nmanageable, chunks and are distributed across many machines (nodes). Data processing functions are distributed on each node and run in parallel, which is referred to as the ",(0,r.kt)("em",{parentName:"p"},"scatter")," phase. When data needs to be aggregated and collected back, it is gathered back into\nindividual nodes, which is referred to as the ",(0,r.kt)("em",{parentName:"p"},"gather")," phase. Spark automatically handles these distribution and collection tasks which frees\nusers from having to worry about complex fault-tolerance and scalability topics. "),(0,r.kt)("p",null,"Spark applications are developed in a variety of programs such as Python, Java, Scala and
1SQL. In this course, we focus on the Python\ninterface for Spark called ",(0,r.kt)("a",{parentName:"p",href:"https://spark.apache.org/docs/latest/api/python/index.html"},"PySpark"),". Spark runs applications through\nits internal interface called SparkContext. The Next sections will explain how to develop and run Spark applications using\nthe SparkContext."),(0,r.kt)("h2",{id:"getting-starting-with-pyspark"},"Getting starting with Pyspark"),(0,r.kt)("p",null,"So how are we going to use this engine? While Spark itself is written in Scala, Pyspark provides us Python programmers a library to run Spark. In order to\ninstall the Pyspark module, run the following to create a python virtualenv and install pyspark:"),(0,r.kt)("pre",null,(0,r.kt)("code",{parentName:"pre",className:"language-bash"},"# create and activate virtualenv if you haven't done so already\npython3.7 -m venv venv\nsource venv/bin/activate\n\n# install pyspark\npip install pyspark\n")),(0,r.kt)("h3",{id:"sparkcontext-how-to-interact-with-data"},"SparkContext: how to interact with data"),(0,r.kt)("p",null,"Now that we have Pyspark installed, let's figure out how to actually use it. Start a new notebook in your Python environment and follow along with these exercises.\nOur first goal is to be able to load and transform our data files. Because Spark is designed to work with data in HDFS (Hadoop Distributed File System) and other\nbig data storage tools as well as local data, it has its own set of data interface objects. A\n",(0,r.kt)("a",{parentName:"p",href:"https://spark.apache.org/docs/latest/api/python/reference/api/pyspark.SparkContext.html#pyspark.SparkContext"},(0,r.kt)("inlineCode",{parentName:"a"},"SparkContext"))," is the most basic of these for connecting to your data. In\nlarger scale operations, a ",(0,r.kt)("inlineCode",{parentName:"p"},"SparkContext")," could be configured to work across clusters or other cloud computing setups but for now, we're just using the little\nlocal cluster that is the machine that is running this notebook. To do this, we create a SparkContext object with the ",(0,r.kt)("inlineCode",{parentName:"p"},"'local[*]'")," options. The ",(0,r.kt)("inlineCode",{parentName:"p"},"[*]")," allows\nthis SparkContext to have access to all local cores; you can manually set this number lower if you would like to limit the context. Let's initialize a context:"),(0,r.kt)("pre",null,(0,r.kt)("code",{parentName:"pre",className:"language-python"},"import pyspark\nsc = pyspark.SparkContext('local[*]')\n")),(0,r.kt)("h3",{id:"resilient-distributed-datasets-rdd"},"Resilient Distributed Datasets (RDD)"),(0,r.kt)("p",null,(0,r.kt)("a",{parentName:"p",href:"https://spark.apache.org/docs/latest/rdd-programming-guide.html"},"Resilient Distributed Dataset or RDD"),'. As the name "distributed" suggests, Spark\nspreads the data out and works on it in parallel. The resiliency comes from the redundant nature of the distributes similar to what was seen in HDFS. In our\nlittle local example, we won\'t be leveraging all of this power, but we can start by taking a simple Python data structure and parallelizing it into a RDD\nwith the following code:'),(0,r.kt)("pre",null,(0,r.kt)("code",{parentName:"pre",className:"language-python"},'\nimport pyspark\nsc = pyspark.SparkContext(\'local[*]\')\n\nlist_of_arrivals = [\n ("PDX", 1),\n ("LAX", 5),\n ("DEN", 3),\n ("PDX", 2),\n ("JFK", 9),\n ("DEN", 5),\n ("PDX", 7),\n ("JFK", 10),\n]\narrivals_rdd = sc.parallelize(list_of_arrivals)\nprint(arrivals_rdd)\n')),(0,r.kt)("p",null,"Great! We have an RDD object! But how do we see what is in our RDD? In order to avoid unnecessary computation or memory usage, Spark uses ",(0,r.kt)("strong",{parentName:"p"},"lazy evaluation"),". "),(0,r.kt)("p",null,"Nothing is computed before it is needed. In order to view the data, we need to tell Spark to collect the distributed data and luckily the ",(0,r.kt)("inlineCode",{parentName:"p"},".collect()")," method\ndoes just that. It is worth noting that ",(0,r.kt)("inlineCode",{parentName:"p"},"collect()")," does exactly what it sounds like and collects all the data distributed across the nodes by the RDD. This\nmeans it is not always the best way to view our data once we are using truly big data that cannot all be collected on one node. For those cases, we can use\ntools like ",(0,r.kt)("inlineCode",{parentName:"p"},"count()")," to see how much data we have without collecting that data together:"),(0,r.kt)("pre",null,(0,r.kt)("code",{parentName:"pre",className:"language-python"},'print("This is our RDD", arrivals_rdd.collect())\nprint("It has {} elements".format(arrivals_rdd.count()))\n')),(0,r.kt)("h3",{id:"filter-data"},"Filter data"),(0,r.kt)("p",null,"Now that we've created our RDD, let's go over some basic functions. First, let's look at filtering data based on a condition. If you're not familiar with Python lambda functions,\n",(0,r.kt)("a",{parentName:"p",href:"https://www.w3schools.com/python/python_lambda.asp"},"here is a quick overview")," We will use this function to filter our data. The ",(0,r.kt)("inlineCode",{parentName:"p"},"x")," in our lambda function is one of the elements\nof our RDD. For this RDD, it is the tuple of ",(0,r.kt)("inlineCode",{parentName:"p"},"('aport code', 'arrival count')"),". Let's filter the data to just Portland arrivals:"),(0,r.kt)("pre",null,(0,r.kt)("code",{parentName:"pre",className:"language-python"},'pdx_arrivals = arrivals_rdd.filter(lambda x: x[0] == "PDX")\nprint(pdx_arrivals.collect())\n')),(0,r.kt)("h3",{id:"group-by-key"},"Group by key"),(0,r.kt)("p",null,"In addition to filtering our data, we also may want to pull together the data based off a key. In our tuples, the first element, ",(0,r.kt)("inlineCode",{parentName:"p"},"airport code"),", functions as a key:"),(0,r.kt)("pre",null,(0,r.kt)("code",{parentName:"pre",className:"language-python"},"grouped_arrivals = arrivals_rdd.groupByKey()\nprint(grouped_arrivals.collect())\n")),(0,r.kt)("p",null,"As you can see above, this creates an iterable object for each airport code. This is the lazy evaluation noted earlier. The iterable allows all the arrival counts to\nbe grouped together without having to load them into memory. If we would like to combine the arrival counts in each group we can use the ",(0,r.kt)("inlineCode",{parentName:"p"},"mapValues()")," function. This\nargument of this function is the function to apply to the iterator. In our case, let's get the count per-airport code so we cause use the basic Python ",(0,r.kt)("inlineCode",{parentName:"p"},"sum")," function as the argument: "),(0,r.kt)("pre",null,(0,r.kt)("code",{parentName:"pre",className:"language-python"},"grouped_arrivals_count = grouped_arrivals.mapValues(sum)\nprint(grouped_arrivals_count.collect())\n")),(0,r.kt)("h2",{id:"conclusion"},"Conclusion"),(0,r.kt)("admonition",{type:"danger"},(0,r.kt)("p",{parentName:"admonition"},"The links in the below paragraph need to be updated!")),(0,r.kt)("p",null,"You just learned the basics of using Spark RDD transformations. To continue learning Spark SQL and DataFrames please refer to\nIntroduction to Apache Spark in /docs/ch3/c3-intro episode of our ",(0,r.kt)("a",{parentName:"p",href:"https://www.datastack.academy/"},"Data Engineering Bootcamp"),". This episode will teach you\nhow to use Spark to read, transform, and write data files. Have fun!"))}c.isMDXComponent=!0},4427:(e,a,t)=>{t.d(a,{Z:()=>n});const n=t.p+"assets/images/parham-5e7600927dc689dca8ab4c199e1f66ef.png"}}]);
Line numbers count LF bytes from the start of the resource, as the search results do. Vendor segments are library code the classifier recognised; they are stored but not indexed. Bytes are shown as Latin1 characters, one per byte.