레이블이 Spark Example인 게시물을 표시합니다. 모든 게시물 표시
레이블이 Spark Example인 게시물을 표시합니다. 모든 게시물 표시

[Spark] Apache Spark Examples

[Spark] Apache Spark Examples

Apache Spark Examples

from : http://spark.apache.org/examples.html

이 예제들은 Spark API의 전체를 간략하게 보여준다. Spark는 분산 데이터셋의 개념으로 만들어 져있다. 이는 Java, Python 객체들을 포함하고 있다. 외부 데이터로 부터 데이터셋을 생성하고 병렬로 그것들을 처리할 수 있다. Spark API의 building block 을 확인해보자. RDD API에는 2가지 타입의 오퍼레이션을 제공하고 있다. transformation으로 이전것으로 부터 새로운 데이터셋을 만들어 내는 작업을 한다. 다음으로 actions로 클러스터 상에서 작업을 시작하는 것이다. Spark의 RDD API의 최상위에는 고수준의 API를 제공하며 DataFrame API와 Machine Learning API가 그것이다. 이 고수준 API는 특정 데이터 오퍼레이션에 대한 간단한 방법으로 작업을 수행할 수 있도록 제공한다. 이 페이지에서는 RDD 고수준 API를 이용한 간단한 예제를 살펴볼 것이다.

RDD API Examples

Word Count

이 예제는 몇가지 transformation으로 데이터셋을 만들고 난 후, (String, Int)쌍의 데이터셋을 만들어 내는 작업을 한다. 이는 count라고 부른다. 그리고 다음으로 결과를 파일로 저장한다.

1. Python 버젼 : 


text_file = sc.textFile("hdfs://...")
counts = text_file.flatMap(lambda line: line.split(" ")) \
.map(lambda word: (word, 1)) \
.reduceByKey(lambda a, b : a + b)
counts.saveAsTextFile("hdfs://...")

2. Scala 버젼 : 

val textFile = sc.textFile("hdfs://...")
val counts = textFile.flatMap(line => line.split(" "))
        .map(word => (word, 1))
        .reduceByKey(_ + _)
counts.saveAsTextFile("hdfs://...")

3. Java 버젼 : 

JavaRDD<String> textFile = sc.textFile("hdfs://...");
JavaRDD<String> words = textFile.flatMap(new FlatMapFunction<String, String>() {
    public Iterable<String> call(String s) { return Arrays.asList(s.split(" ")); }
});
JavaPairRDD<String, Integer> pairs = words.mapToPair(new PairFunction<String, String, Integer>() {
    public Tuple2<String, Integer> call(String s) { return new Tuple2<String, String>(s, 1); }
});
JavaPairRDD<String, Integer> counts = pairs.reduceByKey(new Function2<Integer, Integer, Integer>(){
    public Integer call(Integer a, Integer b) {return a + b;}
});
counts.saveAsTextFile("hdfs://...");

Pi Estimation

Spark는 계산에 집중적인 태스크들에 사용된다. 이 코드는 원에 "throwing darts"를 수행해서 파이를 계산하는 예제이다. 우리는 ((0, 0) to (1, 1)) 사각 정방영역에 랜덤 포인트를 찍고, 얼마나 많이 단위 단위 원 내에 들어가는지 검사하고자 한다. fraction은 Pi / 4이며, 우리는 이 공식을 우리의 계산에 이용할 것이다. 

1. Python 버전 : 

def sample(p) :
    x, y = random(), random()
    return 1 if x*x + y*y < 1 else 0
count = sc.parallelize(xrange(0, NUM_SAMPLES)).map(sample) \
        .reduce(lambda a, b: a + b)
print "Pi is roughly %f" % (4.0 * count / NUM_SAMPLES)

2. Scala 버젼 : 

val count = sc.parallelize(1 to NUM_SAMPLES).map{ i =>
    val x = Math.random()
    val y = Math.random()
    if (x*x + y*y < 1) 1 else 0
}.reduce(_ + _)
println("Pi is roughly " + 4.0 * count / NUM_SAMPLES)

3. Java 버젼 :

List<Integer> l = new ArrayList<Integer>(NUM_SAMPLES);
for (int i = 0; i < NUM_SAMPLES; i++) {
    l.add(i);
}
long count = sc.parallelize(l).filter(new Function<Integer, Boolean>() {
    public Boolean call(Integer i) {
        double x = Math.random();
        double y = Math.random();
        return x * x + y * y < 1;
    }
}).count();
System.out.println("Pi is roughly " + 4.0 * count / NUM_SAMPLES);

DataFrame API Examples

Spark에서 DataFrame는 이름으로 된 칼럼 조합의 분산 데이터 집합이다. 사용자는 DataFrame API를 이용하여 외부 데이터 소스와 스파크의 내장 분산 컬렉션에 대해서 다양한 관계 연산을 수행할 수 있다. 이때 데이터 처리를 위한 별도의 프로시저 없이 이러한 연산이 가능하다. 또한 DataFrame API기반의 프로그램은 자동적으로 스파크 내장 옵티마이저에 의해서 최적화 된다. 

Text Search

이 예제에서는 로그파일에서 에러 메시지를 찾는 예제이다. 

1. Python 버젼 

textFile = sc.textFile("hdfs://...")
# Creates a DataFrame having a single column named "line"
df = textFile.map(lambda r : Row(r)).toDF(["line"])
errors = df.filter(col("line").like("%ERROR%"))
# Counts all the errors
errors.count()
# Counts errors mentioning MySQL
errors.filter(col("line").like("%MySQL%")).count()
# Fetches the MySQL errors as an array of strings
errors.filter(col("line").like("%MySQL%")).count()

2. Scala 버젼 

val textFile = sc.textFile("hdfs://...")
// Creates a DataFrame having a single column named "line"
val df = textFile.toDF("line")
val errors = df.filter(col("line").like("%ERROR%"))
// Counts all the errors
errors.count()
// Counts errors mentioning MySQL
errors.filter(col("line").like("%MySQL%")).count()
// Fetches the MySQL errors as an array of Strings
errors.filter(col("line").like("%MySQL%")).collect()

3. Java 버젼

// Creates a DataFrame having a single column named "line"
JavaRDD<String> textFile = sc.textFile("hdfs://...");
JavaRDD<Row> rowRDD = textFile.map(
  new Function<String, Row>() {
    public Row call(String line) throws Exception {
      return RowFactory.create(line);
    }
  });
List<StructField> fields = new ArrayList<StructField>();
fields.add(DataTypes.createStructField("line", DataTypes.StringType, true));
StructType schema = DataTypes.createStructType(fields);
DataFrame df = sqlContext.createDataFrame(rowRDD, schema);
DataFrame errors = df.filter(col("line").like("%ERROR%"));
// Counts all the errors
errors.count();
// Counts errors mentioning MySQL
errors.filter(col("line").like("%MySQL%")).count();
// Fetches the MySQL errors as an array of strings
errors.filter(col("line").like("%MySQL%")).collect();

Simple Data Operations

이번 예제는 데이터베이스에 있는 테이블을 읽고 각 나이별 사람의 수를 계산하는 예제이다. 최종적으로 우리는 계산된 결과를 JSON타입으로 S3에 저장한다. 단순한 MySQL테이블 "people"에는 "name"과 "age" 2개의 칼럼이 있다. 

1. Python 버젼

# Creates a DataFrame based on a table named "people"
# stored in a MySQL database
url = \
    "jdbc:mysql://yourIP:yourPort/test?user=yourUsername;password=yourPassword"
df = sqlContext \
    .read \
    .format("jdbc") \
    .option("url", url) \
    .option("dbtable", "people") \
    .load()
# Looks the schema of this DataFrame
df.printSchema()
# Counts people by age
countsByAge = df.groupBy("age").count()
countsByAge.show()
# Save countsByAge to S3 in the JSON format.
countsByAge.write.format("json").save("s3a://...")

2. Scala 버젼

// Creates a DataFrame based on a table named "people"
// stored in a MySQL database.
val url =
  "jdbc:mysql://yourIP:yourPort/test?user=yourUsername;password=yourPassword"
val df = sqlContext
  .read
  .format("jdbc")
  .option("url", url)
  .option("dbtable", "people")
  .load()
// Looks the schema of this DataFrame.
df.printSchema()
// Counts people by age
val countsByAge = df.groupBy("age").count()
countsByAge.show()
// Saves countsByAge to S3 in the JSON format.
countsByAge.write.format("json").save("s3a://...")

3. Java 버젼 

// Creates a DataFrame based on a table named "people"
// stored in a MySQL database.
String url =
  "jdbc:mysql://yourIP:yourPort/test?user=yourUsername;password=yourPassword";
DataFrame df = sqlContext
  .read()
  .format("jdbc")
  .option("url", url)
  .option("dbtable", "people")
  .load();
// Looks the schema of this DataFrame.
df.printSchema();
// Counts people by age
DataFrame countsByAge = df.groupBy("age").count();
countsByAge.show();
// Saves countsByAge to S3 in the JSON format.
countsByAge.write().format("json").save("s3a://...");

Machine Learning Example

MLib는 Spark 머신러닝 라이브러리이다. 이는 많은 ML알고리즘을 제공해주고 있다. 이 알고리즘들은 extraction, classification, regression, clustering, recommendation, 등을 제공한다. MLlib 는 또한 워크플로우 구성을 위한 파이프라인을 제공하며, 튜닝 파라미터를 위한 CrossValidator, 그리고 모델의 저장 및 로드를 위한 퍼시스턴스 모델을 제공한다.

Prediction with Learning Example

이 예제에서는 labels과 feature 벡터들의 데이터셋을 획득하고, feature 벡터들에서 부터 labels를 예측하는 학습 수행할 것이다. 이는 Logistic Regression 알고리즘을 이용한다. 

1. Python 버젼

# Every record of this DataFrame contains the label and
# features represented by a vector.
df = sqlContext.createDataFrame(data, ["label", "features"])
# Set parameters for the algorithm.
# Here, we limit the number of iterations to 10.
lr = LogisticRegression(maxIter=10)
# Fit the model to the data.
model = lr.fit(df)
# Given a dataset, predict each point's label, and show the results.
model.transform(df).show()

2. Scala 버젼

// Every record of this DataFrame contains the label and
// features represented by a vector.
val df = sqlContext.createDataFrame(data).toDF("label", "features")
// Set parameters for the algorithm.
// Here, we limit the number of iterations to 10.
val lr = new LogisticRegression().setMaxIter(10)
// Fit the model to the data.
val model = lr.fit(df)
// Inspect the model: get the feature weights.
val weights = model.weights
// Given a dataset, predict each point's label, and show the results.
model.transform(df).show()

3. Java 버젼

// Every record of this DataFrame contains the label and
// features represented by a vector.
StructType schema = new StructType(new StructField[]{
  new StructField("label", DataTypes.DoubleType, false, Metadata.empty()),
  new StructField("features", new VectorUDT(), false, Metadata.empty()),
});
DataFrame df = jsql.createDataFrame(data, schema);
// Set parameters for the algorithm.
// Here, we limit the number of iterations to 10.
LogisticRegression lr = new LogisticRegression().setMaxIter(10);
// Fit the model to the data.
LogisticRegressionModel model = lr.fit(df);
// Inspect the model: get the feature weights.
Vector weights = model.weights();
// Given a dataset, predict each point's label, and show the results.
model.transform(df).show();







[Spark] Quick Start

[Spark] Quick Start

Spark Quick Start

Interactive Analysis with the Spark Shell

기본 : 

스파크 쉘을 다음과 같이 실행하자. 이는 대화형 데이터 분석을 위한 강력한 툴이다.
./bin/pyspark
스파크의 중요한 추상화는 Resilient Distributed Dataset(RDD)라고 부르는 분산 컬렉션이다.
RDD는 Hadoop 입력 포맷 이나 transformation등을 통해서 생성된다. 스파크 패키지에 존재하는 README 파일을 읽어 들이는 예제를 보자.
>>> textFile = sc.textFile("README.md")
RDD는 actions과 transformations를 가진다. action은 반환 값을 가지며, transformation은 새로운 RDD를 생성한다.
다음 몇개의 액션들을 보자.
>>> textFile.count()       # 이 RDD에서 아이템의 총 개수를 센다.
126
>>> textFile.first()          # 이 RDD에서 첫번째 아이템을 출력한다.
u'# Apache Spark'
다음은 transformation을 사용한 예제이다.
이것은 filter transformation을 통해서 새로운 RDD를 생성한다. 이는 file의 subset이다.
>>> linesWithSpark = textFile.filter(lambda line: "Spark" in line)
그리고 transformations과 actions의 체인을 걸어줄 수 있다.
>>> textFile.filter(lambda line: "Spark" in line).count()   # "Spark"를 포함하는 라인이 몇개인가?
15

RDD 연산 더보기 

RDD actions와 transformations는 더 복잡한 계산을 위해서 사용이 가능하다.
가장 긴 단어의 라인을 찾는다고 해보자.
>>> textFile.map(lambda line: len(line.split())).reduce(lambda a, b: a if (a > b) else b)
15
이 첫번째 맵은 라인을 정수 값으로 하는 새로운 RDD를 생성한다. reduce는 가장 긴 라인 카운트를 찾는데 호출된다.map 아규먼트와 recude아규먼트는 (lambda)로 anonymous functions이다. 우리는 최상위 레벨의 파이선 함수를 전달할 수도 있다. 예를 들어 우리는 max함수를 이해하기 쉽게 다음과 같이 만들 수 있다.

>>> def max(a, b) :
...            if a > b:
...                return a
...            else:
...                return b
>>> textFile.map(lambda line: len(line.split())).reduce(max)
15
하나는 Hadoop에 잘 알려진 데이터 흐름 패턴인 MapReduce이다. Spark는 다음과 같이 쉽게 MapReduce를 구현할 수 있다.
>>> wordCounts = textFile.flatMap(lambda line: line.split()).map(lambda word: (word, 1)).reduceByKey(lambda a, b: a+b)
여기에서 우리는 flatMap, map과 reduceByKey transformations를 서로 연결하여 RDD에 존재하는 각 단어의 개수를 새로운 RDD로 (string, int) 쌍으로 만들어 낸다.
우리에 쉘에서는 단어 수를 세는 작업은 collect 액션을 이용하여 구현한다.
>>> wordCounts.collect()
[[(u'and', 9), (u'A', 1), (u'webpage', 1), (u'README', 1), (u'Note', 1), (u'"local"', 1), (u'variable', 1), ...]

Caching

Spark는 또한 여러 클러스들에 넓게 메모리 캐시에 데이터를 넣을 수 있는 기능을 제공한다. 이것은 데이터를 반복적으로 접근할때 매우 유용한 기술이다. 작은 "hot" 데이터셋을 쿼리 할때나, PageRank와 같은 알고리즘을 반복적으로 수행하고자 할때 유용하다. 다음은 linesWithSpark데이터 셋을 캐시 하는 예이다.

>>> linesWithSpark.cache()
>>> linesWithSpark.count()
19
>>> linesWithSpark.count()
19

만약 100라인짜리 텍스트 파일을 캐시한다면 매우 우스운 일일 것이다. 흥미로운 것은 이러한 동일한 함수를 통해서 매우 큰 데이터 셋을 사용할 수 있다는 것이다. 비록 이들이 10개에서 수백개의 노드에 흩어져 있더라도 말이다. 또한 bin/pyspark를 통해서 클러스터에 연결하여 이러한 작업을 수행할 수 있다. 이는 프로그래밍 가이드를 참조하자. https://spark.apache.org/docs/latest/programming-guide.html#initializing-spark

Self-Contained Applications

우리는 Spark API를 이용하여 self-contained application을 작성할 수 있다. 
다음 SimpleApp.py 파일을 살펴보자. 

"""SimpleApp.py"""
from pyspark import SparkContext
logFile = "YOU_SPARK_HOME/README.md"
sc = SparkContext("local", "Simple App")
logData = sc.textFile(logFile).cache()
numAs = logData.filter(lambda s: 'a' in s).count()
numBs = logData.filter(lambda s: 'b' in s).count()
print("Lines with a: %i, lines with b: %i" % (numAs, numBs))
이 프로그램은 텍스트 파일에서 각 라인중에 a와 b를 포함한 라인의 수를 카운트 한다. YOUR_SPARK_HOME은 인스톨된 스파크의 위치를 지정하면 된다. SparkContext는 RDD를 생성하기 위해서 사용된다. 또한 Spark에 Python 함수를 전달할수 있다. 이것은 자동으로 해당 함수를 참조하도록 시리얼라이즈 된다. 어플리케이션에서 사용하기 위해 커스텀 클래스혹은 서드파티 라이브러리를 이용할 수 있다. 또한 이 외부 어플리케이션을 실행하기 위해서 spark-submit를 호출할 수 있다. 자세한 내용은 spark-submit --help를 통해서 사용법을 확인해보자. 

다음 과 같이 함수를 실행해보자. 
# spark-submit 를 실행하여 어플리케이션을 실행한다. 
$ YOUR_SPARK_HOME/bin/spark-submit \
   --master local[4] \
   SimpleApp.py
...
Lines with a: 46, Lines with b: 23