(This page has no text content)
Learning PySpark Build data-intensive applications locally and deploy at scale using the combined powers of Python and Spark 2.0 Tomasz Drabas Denny Lee BIRMINGHAM - MUMBAI
Learning PySpark Copyright © 2017 Packt Publishing All rights reserved. No part of this book may be reproduced, stored in a retrieval system, or transmitted in any form or by any means, without the prior written permission of the publisher, except in the case of brief quotations embedded in critical articles or reviews. Every effort has been made in the preparation of this book to ensure the accuracy of the information presented. However, the information contained in this book is sold without warranty, either express or implied. Neither the authors, nor Packt Publishing, and its dealers and distributors will be held liable for any damages caused or alleged to be caused directly or indirectly by this book. Packt Publishing has endeavored to provide trademark information about all of the companies and products mentioned in this book by the appropriate use of capitals. However, Packt Publishing cannot guarantee the accuracy of this information. First published: February 2017 Production reference: 1220217 Published by Packt Publishing Ltd. Livery Place 35 Livery Street Birmingham B3 2PB, UK. ISBN 978-1-78646-370-8 www.packtpub.com
Credits Authors Tomasz Drabas Denny Lee Reviewer Holden Karau Commissioning Editor Amey Varangaonkar Acquisition Editor Prachi Bisht Content Development Editor Amrita Noronha Technical Editor Akash Patel Copy Editor Safis Editing Project Coordinator Shweta H Birwatkar Proofreader Safis Editing Indexer Aishwarya Gangawane Graphics Disha Haria Production Coordinator Aparna Bhagat Cover Work Aparna Bhagat
(This page has no text content)
Foreword Thank you for choosing this book to start your PySpark adventures, I hope you are as excited as I am. When Denny Lee first told me about this new book I was delighted-one of the most important things that makes Apache Spark such a wonderful platform, is supporting both the Java/Scala/JVM worlds and Python (and more recently R) worlds. Many of the previous books for Spark have been focused on either all of the core languages, or primarily focused on JVM languages, so it's great to see PySpark get its chance to shine with a dedicated book from such experienced Spark educators. By supporting both of these different worlds, we are able to more effectively work together as Data Scientists and Data Engineers, while stealing the best ideas from each other's communities. It has been a privilege to have the opportunity to review early versions of this book, which has only increased my excitement for the project. I've had the privilege of being at some of the same conferences and meetups and watching the authors introduce new concepts in the world of Spark to a variety of audiences (from first timers to old hands), and they've done a great job distilling their experience for this book. The experience of the authors shines through with everything from their explanations to the topics covered. Beyond simply introducing PySpark they have also taken the time to look at up and coming packages from the community, such as GraphFrames and TensorFrames. I think the community is one of those often-overlooked components when deciding what tools to use, and Python has a great community and I'm looking forward to you joining the Python Spark community. So, enjoy your adventure; I know you are in good hands with Denny Lee and Tomek Drabas. I truly believe that by having a diverse community of Spark users we will be able to make better tools useful for everyone, so I hope to see you around at one of the conferences, meetups, or mailing lists soon :) Holden Karau P.S. I owe Denny a beer; if you want to buy him a Bud Light lime (or lime-a-rita) for me I'd be much obliged (although he might not be quite as amused as I am).
About the Authors Tomasz Drabas is a Data Scientist working for Microsoft and currently residing in the Seattle area. He has over 13 years of experience in data analytics and data science in numerous fields: advanced technology, airlines, telecommunications, finance, and consulting he gained while working on three continents: Europe, Australia, and North America. While in Australia, Tomasz has been working on his PhD in Operations Research with a focus on choice modeling and revenue management applications in the airline industry. At Microsoft, Tomasz works with big data on a daily basis, solving machine learning problems such as anomaly detection, churn prediction, and pattern recognition using Spark. Tomasz has also authored the Practical Data Analysis Cookbook published by Packt Publishing in 2016. I would like to thank my family: Rachel, Skye, and Albert—you are the love of my life and I cherish every day I spend with you! Thank you for always standing by me and for encouraging me to push my career goals further and further. Also, to my family and my in-laws for putting up with me (in general). There are many more people that have influenced me over the years that I would have to write another book to thank them all. You know who you are and I want to thank you from the bottom of my heart! However, I would not have gotten through my PhD if it was not for Czesia Wieruszewska; Czesiu - dziękuję za Twoją pomoc bez której nie rozpocząłbym mojej podróży po Antypodach. Along with Krzys Krzysztoszek, you guys have always believed in me! Thank you!
Denny Lee is a Principal Program Manager at Microsoft for the Azure DocumentDB team—Microsoft's blazing fast, planet-scale managed document store service. He is a hands-on distributed systems and data science engineer with more than 18 years of experience developing Internet-scale infrastructure, data platforms, and predictive analytics systems for both on-premise and cloud environments. He has extensive experience of building greenfield teams as well as turnaround/ change catalyst. Prior to joining the Azure DocumentDB team, Denny worked as a Technology Evangelist at Databricks; he has been working with Apache Spark since 0.5. He was also the Senior Director of Data Sciences Engineering at Concur, and was on the incubation team that built Microsoft's Hadoop on Windows and Azure service (currently known as HDInsight). Denny also has a Masters in Biomedical Informatics from Oregon Health and Sciences University and has architected and implemented powerful data solutions for enterprise healthcare customers for the last 15 years. I would like to thank my wonderful spouse, Hua-Ping, and my awesome daughters, Isabella and Samantha. You are the ones who keep me grounded and help me reach for the stars!
About the Reviewer Holden Karau is transgender Canadian, and an active open source contributor. When not in San Francisco working as a software development engineer at IBM's Spark Technology Center, Holden talks internationally on Spark and holds office hours at coffee shops at home and abroad. Holden is a co-author of numerous books on Spark including High Performance Spark (which she believes is the gift of the season for those with expense accounts) & Learning Spark. Holden is a Spark committer, specializing in PySpark and Machine Learning. Prior to IBM she worked on a variety of distributed, search, and classification problems at Alpine, Databricks, Google, Foursquare, and Amazon. She graduated from the University of Waterloo with a Bachelor of Mathematics in Computer Science. Outside of software she enjoys playing with fire, welding, scooters, poutine, and dancing.
www.PacktPub.com eBooks, discount offers, and more Did you know that Packt offers eBook versions of every book published, with PDF and ePub files available? You can upgrade to the eBook version at www.PacktPub. com and as a print book customer, you are entitled to a discount on the eBook copy. Get in touch with us at customercare@packtpub.com for more details. At www.PacktPub.com, you can also read a collection of free technical articles, sign up for a range of free newsletters and receive exclusive discounts and offers on Packt books and eBooks. https://www.packtpub.com/mapt Get the most in-demand software skills with Mapt. Mapt gives you full access to all Packt books and video courses, as well as industry-leading tools to help you plan your personal development and advance your career. Why subscribe? • Fully searchable across every book published by Packt • Copy and paste, print, and bookmark content • On demand and accessible via a web browser
Customer Feedback Thanks for purchasing this Packt book. At Packt, quality is at the heart of our editorial process. To help us improve, please leave us an honest review on this book's Amazon page at https://www.amazon.com/dp/1786463709. If you'd like to join our team of regular reviewers, you can email us at customerreviews@packtpub.com. We award our regular reviewers with free eBooks and videos in exchange for their valuable feedback. Help us be relentless in improving our products!
[ i ] Table of Contents Preface vii Chapter 1: Understanding Spark 1 What is Apache Spark? 2 Spark Jobs and APIs 3 Execution process 3 Resilient Distributed Dataset 4 DataFrames 6 Datasets 6 Catalyst Optimizer 6 Project Tungsten 7 Spark 2.0 architecture 8 Unifying Datasets and DataFrames 9 Introducing SparkSession 10 Tungsten phase 2 11 Structured streaming 13 Continuous applications 14 Summary 15 Chapter 2: Resilient Distributed Datasets 17 Internal workings of an RDD 17 Creating RDDs 18 Schema 20 Reading from files 20 Lambda expressions 21 Global versus local scope 23 Transformations 24 The .map(...) transformation 24 The .filter(...) transformation 25 The .flatMap(...) transformation 26
Table of Contents [ ii ] The .distinct(...) transformation 26 The .sample(...) transformation 27 The .leftOuterJoin(...) transformation 27 The .repartition(...) transformation 28 Actions 29 The .take(...) method 29 The .collect(...) method 29 The .reduce(...) method 29 The .count(...) method 31 The .saveAsTextFile(...) method 31 The .foreach(...) method 32 Summary 32 Chapter 3: DataFrames 33 Python to RDD communications 34 Catalyst Optimizer refresh 35 Speeding up PySpark with DataFrames 36 Creating DataFrames 38 Generating our own JSON data 38 Creating a DataFrame 39 Creating a temporary table 39 Simple DataFrame queries 42 DataFrame API query 42 SQL query 42 Interoperating with RDDs 43 Inferring the schema using reflection 43 Programmatically specifying the schema 44 Querying with the DataFrame API 46 Number of rows 46 Running filter statements 46 Querying with SQL 47 Number of rows 47 Running filter statements using the where Clauses 48 DataFrame scenario – on-time flight performance 49 Preparing the source datasets 50 Joining flight performance and airports 50 Visualizing our flight-performance data 52 Spark Dataset API 53 Summary 54
Table of Contents [ iii ] Chapter 4: Prepare Data for Modeling 55 Checking for duplicates, missing observations, and outliers 56 Duplicates 56 Missing observations 60 Outliers 64 Getting familiar with your data 66 Descriptive statistics 67 Correlations 70 Visualization 71 Histograms 72 Interactions between features 76 Summary 77 Chapter 5: Introducing MLlib 79 Overview of the package 80 Loading and transforming the data 80 Getting to know your data 85 Descriptive statistics 85 Correlations 87 Statistical testing 89 Creating the final dataset 90 Creating an RDD of LabeledPoints 90 Splitting into training and testing 91 Predicting infant survival 91 Logistic regression in MLlib 91 Selecting only the most predictable features 93 Random forest in MLlib 94 Summary 96 Chapter 6: Introducing the ML Package 97 Overview of the package 97 Transformer 98 Estimators 101 Classification 101 Regression 103 Clustering 103 Pipeline 104 Predicting the chances of infant survival with ML 105 Loading the data 105 Creating transformers 106 Creating an estimator 107 Creating a pipeline 107
Table of Contents [ iv ] Fitting the model 108 Evaluating the performance of the model 109 Saving the model 110 Parameter hyper-tuning 111 Grid search 111 Train-validation splitting 115 Other features of PySpark ML in action 116 Feature extraction 116 NLP - related feature extractors 116 Discretizing continuous variables 119 Standardizing continuous variables 120 Classification 122 Clustering 123 Finding clusters in the births dataset 124 Topic mining 124 Regression 127 Summary 129 Chapter 7: GraphFrames 131 Introducing GraphFrames 134 Installing GraphFrames 134 Creating a library 135 Preparing your flights dataset 138 Building the graph 140 Executing simple queries 141 Determining the number of airports and trips 142 Determining the longest delay in this dataset 142 Determining the number of delayed versus on-time/early flights 142 What flights departing Seattle are most likely to have significant delays? 143 What states tend to have significant delays departing from Seattle? 144 Understanding vertex degrees 145 Determining the top transfer airports 146 Understanding motifs 147 Determining airport ranking using PageRank 149 Determining the most popular non-stop flights 151 Using Breadth-First Search 152 Visualizing flights using D3 154 Summary 155
Table of Contents [ v ] Chapter 8: TensorFrames 157 What is Deep Learning? 157 The need for neural networks and Deep Learning 161 What is feature engineering? 163 Bridging the data and algorithm 164 What is TensorFlow? 166 Installing Pip 168 Installing TensorFlow 169 Matrix multiplication using constants 170 Matrix multiplication using placeholders 171 Running the model 172 Running another model 172 Discussion 173 Introducing TensorFrames 174 TensorFrames – quick start 175 Configuration and setup 176 Launching a Spark cluster 176 Creating a TensorFrames library 176 Installing TensorFlow on your cluster 176 Using TensorFlow to add a constant to an existing column 177 Executing the Tensor graph 178 Blockwise reducing operations example 179 Building a DataFrame of vectors 180 Analysing the DataFrame 180 Computing elementwise sum and min of all vectors 181 Summary 182 Chapter 9: Polyglot Persistence with Blaze 183 Installing Blaze 184 Polyglot persistence 185 Abstracting data 186 Working with NumPy arrays 186 Working with pandas' DataFrame 188 Working with files 189 Working with databases 192 Interacting with relational databases 192 Interacting with the MongoDB database 194 Data operations 194 Accessing columns 194 Symbolic transformations 195 Operations on columns 197
Table of Contents [ vi ] Reducing data 198 Joins 200 Summary 202 Chapter 10: Structured Streaming 203 What is Spark Streaming? 203 Why do we need Spark Streaming? 206 What is the Spark Streaming application data flow? 207 Simple streaming application using DStreams 208 A quick primer on global aggregations 213 Introducing Structured Streaming 218 Summary 222 Chapter 11: Packaging Spark Applications 223 The spark-submit command 223 Command line parameters 224 Deploying the app programmatically 227 Configuring your SparkSession 227 Creating SparkSession 228 Modularizing code 229 Structure of the module 229 Calculating the distance between two points 231 Converting distance units 231 Building an egg 232 User defined functions in Spark 232 Submitting a job 233 Monitoring execution 236 Databricks Jobs 237 Summary 241 Index 243
[ vii ] Preface It is estimated that in 2013 the whole world produced around 4.4 zettabytes of data; that is, 4.4 billion terabytes! By 2020, we (as the human race) are expected to produce ten times that. With data getting larger literally by the second, and given the growing appetite for making sense out of it, in 2004 Google employees Jeffrey Dean and Sanjay Ghemawat published the seminal paper MapReduce: Simplified Data Processing on Large Clusters. Since then, technologies leveraging the concept started growing very quickly with Apache Hadoop initially being the most popular. It ultimately created a Hadoop ecosystem that included abstraction layers such as Pig, Hive, and Mahout – all leveraging this simple concept of map and reduce. However, even though capable of chewing through petabytes of data daily, MapReduce is a fairly restricted programming framework. Also, most of the tasks require reading and writing to disk. Seeing these drawbacks, in 2009 Matei Zaharia started working on Spark as part of his PhD. Spark was first released in 2012. Even though Spark is based on the same MapReduce concept, its advanced ways of dealing with data and organizing tasks make it 100x faster than Hadoop (for in-memory computations). In this book, we will guide you through the latest incarnation of Apache Spark using Python. We will show you how to read structured and unstructured data, how to use some fundamental data types available in PySpark, build machine learning models, operate on graphs, read streaming data, and deploy your models in the cloud. Each chapter will tackle different problem, and by the end of the book we hope you will be knowledgeable enough to solve other problems we did not have space to cover here.
Preface [ viii ] What this book covers Chapter 1, Understanding Spark, provides an introduction into the Spark world with an overview of the technology and the jobs organization concepts. Chapter 2, Resilient Distributed Datasets, covers RDDs, the fundamental, schema-less data structure available in PySpark. Chapter 3, DataFrames, provides a detailed overview of a data structure that bridges the gap between Scala and Python in terms of efficiency. Chapter 4, Prepare Data for Modeling, guides the reader through the process of cleaning up and transforming data in the Spark environment. Chapter 5, Introducing MLlib, introduces the machine learning library that works on RDDs and reviews the most useful machine learning models. Chapter 6, Introducing the ML Package, covers the current mainstream machine learning library and provides an overview of all the models currently available. Chapter 7, GraphFrames, will guide you through the new structure that makes solving problems with graphs easy. Chapter 8, TensorFrames, introduces the bridge between Spark and the Deep Learning world of TensorFlow. Chapter 9, Polyglot Persistence with Blaze, describes how Blaze can be paired with Spark for even easier abstraction of data from various sources. Chapter 10, Structured Streaming, provides an overview of streaming tools available in PySpark. Chapter 11, Packaging Spark Applications, will guide you through the steps of modularizing your code and submitting it for execution to Spark through command-line interface. For more information, we have provided two bonus chapters as follows: Installing Spark: https://www.packtpub.com/sites/default/files/downloads/ InstallingSpark.pdf Free Spark Cloud Offering: https://www.packtpub.com/sites/default/files/ downloads/FreeSparkCloudOffering.pdf
Preface [ ix ] What you need for this book For this book you need a personal computer (can be either Windows machine, Mac, or Linux). To run Apache Spark, you will need Java 7+ and an installed and configured Python 2.6+ or 3.4+ environment; we use the Anaconda distribution of Python in version 3.5, which can be downloaded from https://www.continuum. io/downloads. The Python modules we randomly use throughout the book come preinstalled with Anaconda. We also use GraphFrames and TensorFrames that can be loaded dynamically while starting a Spark instance: to load these you just need an Internet connection. It is fine if some of those modules are not currently installed on your machine – we will guide you through the installation process. Who this book is for This book is for everyone who wants to learn the fastest-growing technology in big data: Apache Spark. We hope that even the more advanced practitioners from the field of data science can find some of the examples refreshing and the more advanced topics interesting. Conventions In this book, you will find a number of styles of text that distinguish between different kinds of information. Here are some examples of these styles, and an explanation of their meaning. Code words in text, database table names, folder names, filenames, file extensions, pathnames, dummy URLs, user input, and Twitter handles are shown as follows: A block of code is set as follows: data = sc.parallelize( [('Amber', 22), ('Alfred', 23), ('Skye',4), ('Albert', 12), ('Amber', 9)]) When we wish to draw your attention to a particular part of a code block, the relevant lines or items are set in bold: rdd1 = sc.parallelize([('a', 1), ('b', 4), ('c',10)]) rdd2 = sc.parallelize([('a', 4), ('a', 1), ('b', '6'), ('d', 15)]) rdd3 = rdd1.leftOuterJoin(rdd2)
Loading comments...
Reply to Comment
Edit Comment