Thursday, 20 August 2015

Learning Spark - Chapter 8 Notes

Configuration precedence for Spark

1) 1st Priority to configuration used with set() in code
2) 2nd Priority to Command line arguments
3) 3rd Priority to Properties File
4) 4th Priority to Default Properties.




Configuration Property
Default Value
Description
spark.executor.memory
(--executor-memory)
-512m
Amount of memory to use per executor process, in the same format as JVM memory strings (e.g., 512m, 2g)
spark.executor.cores

(--executor-cores
--totalexecutor-cores)
1
Configurations for bounding the number of cores used by the application. In YARN mode spark.executor.cores will assign a specific number of cores to each executor. In standalone and Mesos modes, you can upper-bound the total number of cores across all executors using spark.cores.max.
spark.speculation
False
Setting to true will enable speculative execution of tasks. This means tasks that are running slowly will have a second copy launched on another node. Enabling this can help cut down on straggler tasks in large clusters.
spark.storage.
blockManager
TimeoutIntervalMs
45000
An internal timeout used for tracking the liveness of executors. For jobs that have long garbage collection pauses, tuning this to be 100 seconds (a value of 100000) or higher can prevent thrashing. In future versions of Spark this may be replaced with a general timeout setting, so check current documentation.

spark.executor.
extraJavaOptions                           

spark.executor.
extraClassPath

spark.executor.
extraLibraryPath


These three options allow you to customize the launch behavior of executor JVMs. The three flags add extra Java options, classpath entries, or path entries for the JVM library path. These parameters should be specified as strings (e.g., spark.executor.extraJavaOptions="- XX:+PrintGCDetails-XX:+PrintGCTi meStamps"). Note that while this allows you to manually augment the executor classpath, the recommended way to add dependencies is through the --jars flag to spark-submit (not using this option)
spark.serializer
org.apache.spark.serializer.JavaSerializer
Class to use for serializing objects that will be sent over the network or need to be cached in serialized form. The default of Java Serialization works with any serializable Java object but is quite slow, so we recommend using org.apache.spark.seri alizer.KryoSerializer and configuring Kryo serialization when speed is necessary. Can be any subclass of org.apache.spark.Serial izer.
spark.eventLog.enabled
False
Set to true to enable event logging, which allows completed Spark jobs to be viewed using a history server.
spark.eventLog.dir
file:///<file_path>          
The storage location used for event logging, if enabled. This can set to any File System eg:HDFS.





SPARK_LOCAL_DIRS     Local Directory that spark should use for Shuffle. This property cannot be set using the configuration object.

Spark's default cache operation persists using MEMORY_ONLY storage level.

MEMORY_AND_DESK persists RDDs to Memory and if the storage is not sufficient, then drops the old RDDs to DISK and reads them back to memory, when they are needed.


Some Key points about Spark Shuffle Operation

What is Shuffle.?

The shuffle is Spark’s mechanism for re-distributing data so that is grouped differently across partitions. This typically involves copying data across executors and machines, making the shuffle a complex and costly operation.


Performance Impact

  • The Shuffle is an expensive operation since it involves disk I/O, data serialization, and network I/O. To organize data for the shuffle, Spark generates sets of tasks - map tasks to organize the data, and a set of reduce tasks to aggregate it. 
  •  Certain shuffle operations can consume significant amounts of heap memory since they employ in-memory data structures to organize records before or after transferring them. Specifically, reduceByKey and aggregateByKey create these structures on the map side and 'ByKey operations generate these on the reduce side. When data does not fit in memory Spark will spill these tables to disk, incurring the additional overhead of disk I/O and increased garbage collection.
  •  Shuffle also generates a large number of intermediate files on disk. As of Spark 1.3, these files are not cleaned up from Spark’s temporary storage until Spark is stopped, which means that long-running Spark jobs may consume available disk space. This is done so the shuffle doesn’t need to be re-computed if the lineage is re-computed. The temporary storage directory is specified by the spark.local.dir configuration parameter when configuring the Spark context.
  • Operations that involve Shuffle 
                      repartition, coalesce, By Key Operations Except Count By Key, Join, Co Group

  
ReduceByKey vs. GroupByKey

Reduce By Key is more efficient than Group By Key, as Reduce involves grouping operations at each partition. A concept similar to Hadoop's Combiner


Word Count example to understand the difference between Reduce By Key and Group By Key.



Reduce By Key:  The reduce operation happens for each partition, and then the results will be shuffled.



Group By Key: Each and every element gets suffled first, and then grouping operation happens. 

Conclusion: Reduce By Key might not fit or solve all the Group By Operations, but wherever possible, it is better to opt for Reduce By Key, Fold By Key or Combine By Key, instead of Group By.

Sunday, 16 August 2015

Learning Spark - Chapter 7 Notes

Spark has pluggable cluster manager. We can configure it with Yarn, Mesos Cluster and can work with Spark's Stand Alone Spark Cluster.


In distributed mode, Spark uses a master/slave architecture with one central coordinator

and many distributed workers. The central coordinator is called the driver.


The driver communicates with a potentially large number of distributed workers called

executors. The driver runs in its own Java process and each executor is a separate Java

process. A driver and its executors are together termed a Spark application.

Driver is process where main() method of your program runs and it is the process running user code that creates Spark Context, RDDs and performs transformations and actions.




When the driver runs, it performs two duties:



1) It is responsible for converting a user program into units of physical execution called tasks.

 At a high level, all Spark programs follow the same structure:

    They create RDDs from some input, derive new RDDs from those using transformations, and perform actions to collect or save data. A Spark program implicitly creates a logical directed acyclic graph (DAG) of operations.When the driver runs, it converts this logical graph into a physical execution


Spark performs several optimizations, such as “pipelining” map transformations together to merge them, and converts the execution graph into a set of stages. Each stage, in turn, consists of multiple tasks. The tasks are bundled up and prepared to be sent to the cluster. Tasks are the smallest unit of work in Spark; a typical user program can launch hundreds or thousands of individual tasks.


2)Scheduling tasks on executors


Given a physical execution plan, a Spark driver must coordinate the scheduling of individual tasks on executors. When executors are started they register themselves with the driver, so it has a complete view of the application’s executors at all times. Each executor represents a process capable of running tasks and storing RDD data. The Spark driver will look at the current set of executors and try to schedule each task in an appropriate location, based on data placement. When tasks execute, they may have a side effect of storing cached data. The driver also tracks the location of cached data and uses it to schedule future tasks that access that data.






The driver exposes information about the running Spark application through a web interface, which by default is available at port 4040. For instance, in local mode, this UI is available at http://localhost:4040


Executors
Spark executors are worker processes responsible for running the individual tasks in
a given Spark job. Executors are launched once at the beginning of a Spark application
and typically run for the entire lifetime of an application, though Spark applications
can continue if executors fail. Executors have two roles. First, they run the tasks
that make up the application and return results to the driver. Second, they provide
in-memory storage for RDDs that are cached by user programs, through a service
called the Block Manager that lives within each executor. Because RDDs are cached
directly inside of executors, tasks can run alongside the cached data.


Values for master

spark://host:port -> connect to a Spark Standalone cluster at the specified port. By default Spark Standalone masters use port 7077.
mesos://host:port  -> Connect to a Mesos cluster master at the specified port. By default Mesos masters listen on port 5050.
yarn -> Connect to a YARN cluster. When running on YARN you’ll need to set the HADOOP_CONF_DIR environment variable to point the location of your Hadoop configuration directory, which contains information about the cluster.

local Run in local mode with a single core.
local[N] Run in local mode with N cores.
local[*] Run in local mode and use as many cores as the machine has.

Deploy Mode Importance:

Whether to launch the driver program locally (“client”) or on one of the worker machines inside the cluster (“cluster”). In client mode spark-submit will run your driver on the same machine where spark-submit is itself being invoked. In cluster mode, the driver will be shipped to execute on a worker node in the cluster. The default is client mode. 

Few Points about Mesos:
 
Unlike the other cluster managers, Mesos offers two modes to share resources between executors on the same cluster.

1) Fine Grained Mode: In “fine-grained” mode, which is the default, executors scale up and down the number of CPUs they claim from Mesos as they execute tasks, and so a machine running multiple executors can dynamically share CPU resources between them.

2) “coarse-grained” mode, Spark allocates a fixed number of CPUs to each executor in advance and never releases them until the application ends, even if the executor is not currently running tasks.

You can enable coarse-grained: mode by passing --conf spark.mesos.coarse=true to spark-submit

The fine-grained Mesos mode is attractive when multiple users share a cluster to run interactive workloads such as shells, because applications will scale down their number of cores when they’re not doing work and still allow other users’ programs to use the cluster. The downside, however, is that scheduling tasks through fine-grained mode adds more latency

you can use a mix of scheduling modes in the same Mesos cluster (i.e., some of your Spark applications might have spark.mesos.coarse set to true and some might not).

Spark on Mesos supports running applications only in the “client” deploy mode—that is, with the driver running on the machine that submitted the application. If you would like to run your driver in the Mesos cluster as well, frameworks like Aurora and Chronos allow you to submit arbitrary scripts to run on Mesos and monitor them.

 


Some Good Books and download links

Some Important Books and good reads.

1) Machine Learning: Hands-On for Developers and Technical Professionals
Download Link:
http://www.topitbooks.com/machine-learning-hands-developers-technical-professionals-4473.html

2) A nice article on machine learning.

https://kaggle2.blob.core.windows.net/forum-message-attachments/8888/DZone%20introduction%20to%20ML.pdf?sv=2012-02-12&se=2015-07-19T07%3A22%3A04Z&sr=b&sp=r&sig=rLt0FScMx1FxlCFxYnasZ8mMhCnlemHA4yfHV%2BmTMUk%3D

3) Using Flume:
https://www.geekbooks.me/book/view/using-flume

4) Apche Hadoop Yarn by Arun C Murthy

http://www.topitbooks.com/apache-hadoop-yarn-2512.html

5) Advanced Analytics with Spark

http://richmagbooks.com/advanced-analytics-spark-patterns-learning-data-scale/

6) Important article on LZO Compression
https://altiscale.zendesk.com/hc/en-us/articles/202180153-Compressing-and-Indexing-your-Data-with-LZO

http://stackoverflow.com/questions/23560281/do-we-need-to-create-an-index-file-with-lzop-if-compression-type-is-record-ins

http://blog.cloudera.com/blog/2009/11/hadoop-at-twitter-part-1-splittable-lzo-compression/

 

Scala and Spark Notes.

Adding dependency in SBT:

Eg:

libraryDependencies += "org.apache.spark" %% "spark-core" % "1.4.1";

Sample SBT:

name := "spark_srini_scala"

version := "1.0.0"

scalaVersion := "2.11.5"

libraryDependencies += "org.apache.spark" %% "spark-core" % "1.4.1";
libraryDependencies += "com.fasterxml.jackson.core" % "jackson-databind" % "2.2.2";
libraryDependencies += "com.fasterxml.jackson.module" %% "jackson-module-scala" % "2.2.2";

Adding SBT Eclipse Plugin:
Go to scala folder in your system:
eg: SBT Folder in my system

C:\Users\lenovo\.sbt\0.13

Create a folder with name plugins in this folder.
and create plugins.sbt file with the following lines

addSbtPlugin("com.typesafe.sbteclipse" % "sbteclipse-plugin" % "4.0.0")

Building Scala Project: Command

sbt clean package

Reading and Writing JSON in Scala
1) Create a top Level Class with Person
2) Try below code

val person = Person("fred", 25)
val mapper = new ObjectMapper()
mapper.registerModule(DefaultScalaModule)    

val out = new StringWriter
mapper.writeValue(out, person)
val json = out.toString()
println(json)

val person2 = mapper.readValue(json, classOf[Person])
println(person2)


Reading JSON File in Spark:

package spark

import org.apache.spark.SparkConf
import org.apache.spark.SparkContext

import com.fasterxml.jackson.databind.DeserializationFeature
import com.fasterxml.jackson.databind.ObjectMapper
import com.fasterxml.jackson.module.scala.DefaultScalaModule

/**
 * @author Srini
 */

case class JsonInput(name: String, occupation: String);
 
 object JSONExample { 
 def main(args: Array[String]) {
    val conf = new SparkConf();
    val sc = new SparkContext(conf);

    val lines = sc.textFile(args(0));
    
    println("========= input file is args(0)");
    
    val result = lines.mapPartitions(y => {
      val mapper = new ObjectMapper();
      mapper.registerModule(DefaultScalaModule);
    
      y.flatMap( x => {println("line is " + x);Some(mapper.readValue(x, classOf[JsonInput]))}) 
    }, true)
    
    val mapper = new ObjectMapper();  
    
    val y = result.map(mapper.writeValueAsString(_))
    y.saveAsTextFile(args(1))
  }
}

Saving as Hadoop File:
 
public static class ConvertToWritableTypes implements
PairFunction<Tuple2<String, Integer>, Text, IntWritable> {
            public Tuple2<Text, IntWritable> call(Tuple2<String, Integer> record) {
            return new Tuple2(new Text(record._1), new IntWritable(record._2));
       }
}

JavaPairRDD<String, Integer> rdd = sc.parallelizePairs(input);
JavaPairRDD<Text, IntWritable> result = rdd.mapToPair(new ConvertToWritableTypes());
result.saveAsHadoopFile(fileName, Text.class, IntWritable.class,
SequenceFileOutputFormat.class);  
 
Reading Compressed file can be done using newHadoopApiFile
 
sc.newAPIHadoopFile(inputFile, classOf[LzoJsonInputFormat],
classOf[LongWritable], classOf[MapWritable], conf)
 
 
Writing ProtoBuf file in Spark


  val job = new Job()
  val conf = job.getConfiguration
  LzoProtobufBlockOutputFormat.setClassConf(classOf[Places.Venue], conf);
  val dnaLounge = Places.Venue.newBuilder()
  dnaLounge.setId(1);
  dnaLounge.setName("DNA Lounge")
  dnaLounge.setType(Places.Venue.VenueType.CLUB)
  val data = sc.parallelize(List(dnaLounge.build()))
  val outputData = data.map { pb =>
    val protoWritable = ProtobufWritable.newInstance(classOf[Places.Venue]);
    protoWritable.set(pb)
    (null, protoWritable)
  }
  outputData.saveAsNewAPIHadoopFile(outputFile, classOf[Text],
    classOf[ProtobufWritable[Places.Venue]],
    classOf[LzoProtobufBlockOutputFormat[ProtobufWritable[Places.Venue]]], conf)


Scala's textFile can read the Compressed Files automatically, but it will disable splitting.
If Splitting is important, go for newApiHadoopFile

Accessing Hive data from Spark:

import org.apache.spark.sql.hive.HiveContext
val hiveCtx = new org.apache.spark.sql.hive.HiveContext(sc)
val rows = hiveCtx.sql("SELECT name, age FROM users")
val firstRow = rows.first()
println(firstRow.getString(0))

JSONs can be simply accessed using the Hive Context. with jsonFile function. 

val tweets = hiveCtx.jsonFile("tweets.json")
tweets.registerTempTable("tweets")
val results = hiveCtx.sql("SELECT user.name, text FROM tweets")

JDBC Load Sample
 
 object LoadSimpleJdbc {
      def main(args: Array[String]) {
        if (args.length < 1) {
          println("Usage: [sparkmaster]")
          exit(1)
        }
        val master = args(0)
        val sc = new SparkContext(master, "LoadSimpleJdbc", System.getenv("SPARK_HOME"))
        val data = new JdbcRDD(sc,
          createConnection, "SELECT * FROM panda WHERE ? <= id AND ID <= ?",
          lowerBound = 1, upperBound = 3, numPartitions = 2, mapRow = extractValues)
        println(data.collect().toList)
      }
      def createConnection() = {
        Class.forName("com.mysql.jdbc.Driver").newInstance();
        DriverManager.getConnection("jdbc:mysql://localhost/test?user=holden");
      }
      def extractValues(r: ResultSet) = {
        (r.getInt(1), r.getString(2))
      }
    }


we provide a query that can read a range of the data, as well as a lower
Bound and upperBound value for the parameter to this query. These parameters
allow Spark to query different ranges of the data on different machines, so we
don’t get bottlenecked trying to load all the data on a single node.
 
No SQL DBs: 
For Hbase we need to set below parameter
 
conf.set(TableInputFormat.INPUT_TABLE, "tablename")
 
For Cassandra

conf.set("spark.cassandra.connection.host", "hostname")
Reading from Cassandra can be done using 
sc.cassandraTable

For Elastic Search
 
val jobConf = new JobConf(sc.hadoopConfiguration)
jobConf.set(ConfigurationOptions.ES_RESOURCE_READ, args(1))
jobConf.set(ConfigurationOptions.ES_NODES, args(2)) 
 
For HBase and Elastic Search

Elastic Search:
val currentTweets = sc.hadoopRDD(jobConf,
classOf[EsInputFormat[Object, MapWritable]], classOf[Object],
classOf[MapWritable])

HBase:

val rdd = sc.newAPIHadoopRDD(
conf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result])


Learning Spark Notes: Advanced Spark Programming - Chapter 6

Accumulators:

These are global variables and they can be accessed even in distributed mode.

They need three things

1) Declare an accumulator using sc.accumulator(<initialValue>)
2) if scala, then use += to add any value of your wish (eg: accum += 1), if java, then method add
(eg: accum.add(1))
3) get the value of the accumulator using value() method (eg: accum.value())

Note: Accumulators generally used in transformations. In order get the actual value, you need to perform an action first.  such as, saveAsTextFile, count etc.

Accumulators are write only variables

Accumulators in transformations may be calculated multiple times in transformations if there are failures or side effects. So its always advisable to use them with actions.

Custom Accumulators:

Custom Accumulator has to extend AccumulatorParam  and tow methods
1) for Initialization - method zero
2) for Addition - method addInPlace

Scala Example for Custom Accumulator

object CustomAccumulatorExmp {
  def main(args: Array[String]) {
    val conf = new SparkConf
    val sc = new SparkContext(conf)

    val acc = sc.accumulator(List(0))(ListAccumulator)
    val lines = sc.textFile(args(0));

    val lengths = lines.map { string =>
      {
        acc += List(string.size);
        println("length is " + string.length);
        string.length();
      }
    }

    lengths.count();
    lengths.first()
    acc.value.foreach { x => println(" -- " + x); }

  }

  object ListAccumulator extends AccumulatorParam[List[Int]] {
    def zero(initialValue: List[Int]): List[Int] = {
      return initialValue
    }

    def addInPlace(v1: List[Int], v2: List[Int]): List[Int] = {
      v1 ::: v2
    }
  }
}


BroadCast Variables
BroadCast variables are shared variables across the worker nodes, for one or more operations. They come in handy for Look Up kind of operations, or Map Joins etc.

It is important to chose right Serialization format that is fast and

Example Code in Java:

I have written an example to calculate simple sentiment analysis in java using broadcast variables and accumulators.

Note:This usecase is just to demonstrate accumulators and broadcast variables. So, please do not mind about the use case.

The code written to analyze the sentiment for Two Pro Kabaddi League teams Hyderabad and Banglore from given text. Simply tokenizes the whole text and searches for pre-defined set of words that maps to both the teams and when encountered the match, accumulates the points.

Percentage gets calculated based on the total points in the end.


import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;

import org.apache.spark.Accumulator;
import org.apache.spark.SparkConf;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.JavaSparkContext;
import org.apache.spark.api.java.function.FlatMapFunction;
import org.apache.spark.api.java.function.VoidFunction;
import org.apache.spark.broadcast.Broadcast;

public class SentimentAnalysisWithBoradCast {
   
    public static void main(String... args){
       
        Map<String, Integer> teluguMap = new HashMap<String, Integer>();
        teluguMap.put("Titans",1);
        teluguMap.put("Hyderabad",1);
        teluguMap.put("Telugu",1);
        teluguMap.put("Andhra",1);
        teluguMap.put("Nizam",1);
        teluguMap.put("Rahul",1);
       
        Map<String, Integer> bangloreMap = new HashMap<String, Integer>();
        bangloreMap.put("Bulls",1);
        bangloreMap.put("Blore",1);
        bangloreMap.put("Banglore",1);
        bangloreMap.put("Mohit",1);
        bangloreMap.put("Chillar",1);
        bangloreMap.put("Karnataka",1);
       
        SparkConf conf = new SparkConf();
        JavaSparkContext sc = new JavaSparkContext(conf);
        final Broadcast<Map<String, Integer>> telugu = sc.broadcast(teluguMap);
        final Broadcast<Map<String, Integer>> banglore = sc.broadcast(bangloreMap);
       
        final Accumulator<Integer> teluguPoints= sc.accumulator(0);
        final Accumulator<Integer> banglorePoints= sc.accumulator(0);
       
        JavaRDD<String> lines = sc.textFile(args[0]);
       
        JavaRDD<String> words = lines.flatMap(new FlatMapFunction<String, String>(){
            public Iterable<String> call(String input){
                return Arrays.asList(input.split(" "));
            }
        });
       
        words.foreach(new VoidFunction<String>(){
            public void call(String input){
               
                if(telugu.value().get(input) != null)
                {
                    teluguPoints.add(1);
                }
               
                if(banglore.value().get(input) != null)
                {
                    banglorePoints.add(1);
                }
            }
        });
       
        int totalPoints = banglorePoints.value() + teluguPoints.value();
       
        System.out.printf("percentage of Banglore Bulls Sentiment is " + (banglorePoints.value()/(totalPoints * 1.00)) * 100);
        sc.close(); 
    }
}



Pipes:

Can execute any executable file and pass arguments to it using pipe.
Example

import org.apache.spark.SparkConf
import org.apache.spark.SparkContext
import org.apache.spark.SparkFiles


/**
 * @author lenovo
 */

object PipesExample {
  def main(args:Array[String]){
    val conf = new SparkConf();
    val sc = new SparkContext(conf);
   
    val distScript = "/usr/local/test.pl"
    sc.addFile(distScript)
   
    val rdd = sc.parallelize(Array("37.75889318222431"))
   
    val piped = rdd.pipe(Seq(SparkFiles.get("test.pl")),Map("SEPARATOR" -> ","))
   
    println(" -- " + piped.collect().mkString(" "));
   
  }
}