Apache Spark 빠른 시작
Apache Spark 빠른 시작 (Quick Start)
Spark를 처음 접하는 분들이 가장 먼저 해보면 좋은 입문 튜토리얼이에요. 인터랙티브 셸(PySpark나 스파크 셸)로 API를 하나씩 눌러보는 것부터 시작해서, 자립형 애플리케이션을 Scala·Java·Python으로 작성하고 실행하는 흐름까지 한번에 따라가 볼게요. RDD보다는 Dataset/DataFrame을 쓰는 게 성능에서 유리하다는 현재 공식 권장 방향도 여기서 확실히 잡아둡니다.
시작 전에 준비할 것
공식 사이트에서 배포본을 내려받아 두면 돼요. 이 튜토리얼은 HDFS를 쓰지 않으니, 어떤 Hadoop 버전용 패키지든 상관없이 받으면 됩니다.
Spark 2.0 이전에는 주된 프로그래밍 인터페이스가 RDD였어요. 2.0부터는 RDD를 Dataset으로 대체했는데, RDD처럼 강타입(strongly-typed)이면서 내부적으로 더 풍부한 최적화를 제공합니다. RDD 인터페이스 자체는 여전히 지원되지만, 성능이 더 좋은 Dataset으로 전환하는 걸 공식 문서에서도 강하게 권장하고 있어요.
스파크 셸에서 인터랙티브 분석
기본 (Basics)
셸은 API를 익히는 가장 간단한 방법이면서, 동시에 데이터를 인터랙티브하게 분석하는 강력한 도구예요. Spark 디렉터리에서 다음과 같이 실행합니다.
./bin/pyspark
현재 환경에 PySpark가 pip으로 설치돼 있다면 그냥 pyspark라고 입력해도 돼요.
Spark의 핵심 추상화는 Dataset입니다. Dataset은 구조화된 정보의 집합으로, HDFS 같은 Hadoop InputFormat에서 만들거나 다른 Dataset을 변환해서 만들 수 있어요. Python은 동적 타입이라 구현 수준에서 모든 Dataset이 Dataset[Row]가 돼요. 여기서 또 하나의 핵심 개념인 DataFrame(이름이 붙은 컬럼을 가진 Dataset)이 등장합니다. pandas나 R의 DataFrame에 익숙하다면 Spark DataFrame도 비슷하게 바로 적응할 수 있을 거예요.
Spark 소스 디렉터리의 README.md 파일을 DataFrame으로 만들어 볼게요.
>>> textFile = spark.read.text("README.md")
DataFrame을 만들면 액션을 수행하거나 다른 DataFrame으로 변환할 수 있어요.
>>> textFile.count() # 이 DataFrame의 행 수
126
>>> textFile.first() # 첫 번째 행
Row(value=u'# Apache Spark')
이제 filter 함수로 일부 행만 담긴 새 DataFrame을 만들어 봅니다.
>>> linesWithSpark = textFile.filter(textFile.value.contains("Spark"))
변환과 액션은 이렇게 체이닝할 수도 있어요.
>>> textFile.filter(textFile.value.contains("Spark")).count() # "Spark"가 들어간 행은 몇 줄?
15
Dataset 연산 더 알아보기
Dataset의 액션과 변환은 더 복잡한 계산에도 쓰입니다. 단어가 가장 많은 줄을 찾아볼게요.
>>> from pyspark.sql import functions as sf
>>> textFile.select(sf.size(sf.split(textFile.value, "\s+")).name("numWords")).agg(sf.max(sf.col("numWords"))).collect()
[Row(max(numWords)=15)]
각 줄을 정수 값에 매핑하고 별칭 numWords를 붙여 새 DataFrame을 만든 뒤, agg로 가장 큰 단어 수를 찾는 구조예요. select와 agg의 인자는 모두 Column이며, df.colName으로 컬럼을 가져오거나 pyspark.sql.functions가 제공하는 편리한 함수들로 새 Column을 만듭니다.
Hadoop이 대중화시킨 MapReduce 패턴도 Spark로 쉽게 구현돼요.
>>> wordCounts = textFile.select(sf.explode(sf.split(textFile.value, "\s+")).alias("word")).groupBy("word").count()
select 안의 explode로 줄 단위 Dataset을 단어 단위 Dataset으로 바꾸고, groupBy와 count를 조합해 "word"와 "count" 두 컬럼의 DataFrame으로 단어별 개수를 계산하는 흐름이에요.
>>> wordCounts.collect()
[Row(word=u'online', count=1), Row(word=u'graphs', count=1), ...]
캐싱 (Caching)
Spark는 데이터셋을 클러스터 전체의 인메모리 캐시로 끌어올리는 것도 지원해요. 자주 접근하는 작은 "핫" 데이터셋을 조회하거나 PageRank 같은 반복 알고리즘을 돌릴 때 특히 유용합니다.
>>> linesWithSpark.cache()
>>> linesWithSpark.count()
15
>>> linesWithSpark.count()
15
100줄짜리 텍스트 파일을 탐색하고 캐싱하는 게 우스워 보일 수 있지만, 이 함수들은 수십~수백 개 노드에 걸쳐 쪼개진 아주 큰 데이터셋에도 그대로 적용된다는 게 핵심이에요. bin/pyspark를 클러스터에 연결하면 같은 작업을 인터랙티브하게 할 수도 있습니다.
자립형 애플리케이션 (Self-Contained Applications)
셸이 아니라 스스로 실행되는 애플리케이션을 PySpark로 작성해 볼게요. 예시로 SimpleApp.py를 만듭니다.
"""SimpleApp.py"""
from pyspark.sql import SparkSession
logFile = "YOUR_SPARK_HOME/README.md" # 시스템에 있는 어떤 파일이든
spark = SparkSession.builder.appName("SimpleApp").getOrCreate()
logData = spark.read.text(logFile).cache()
numAs = logData.filter(logData.value.contains('a')).count()
numBs = logData.filter(logData.value.contains('b')).count()
print("Lines with a: %i, lines with b: %i" % (numAs, numBs))
spark.stop()
이 프로그램은 텍스트 파일에서 'a'와 'b'가 들어간 줄 수를 세는 단순한 예시예요. YOUR_SPARK_HOME은 실제 Spark 설치 경로로 바꿔야 합니다. 셸과 달리 애플리케이션에서는 SparkSession을 프로그램 안에서 직접 초기화하는데, SparkSession.builder로 만들고 애플리케이션 이름을 지정한 뒤 getOrCreate로 인스턴스를 얻는 구조예요.
커스텀 클래스나 서드파티 라이브러리를 쓰는 경우, 코드 의존성을 .zip으로 묶어 spark-submit의 --py-files 인자로 추가할 수도 있어요 (자세한 건 spark-submit --help 참고).
bin/spark-submit 스크립트로 애플리케이션을 실행합니다.
# spark-submit으로 실행
$ YOUR_SPARK_HOME/bin/spark-submit \
--master "local[4]" \
SimpleApp.py
...
Lines with a: 46, Lines with b: 23
PySpark가 pip으로 설치돼 있다면 일반 Python 인터프리터로도 실행할 수 있어요.
$ python SimpleApp.py
...
Lines with a: 46, Lines with b: 23
Scala 애플리케이션이라면 sbt 설정 파일 build.sbt에 Spark를 의존성으로 선언하고, Java라면 Maven pom.xml에 spark-sql_{{site.SCALA_BINARY_VERSION}} 아티팩트를 추가한 뒤 mvn package로 JAR를 만들어 --class와 함께 spark-submit로 실행하면 됩니다. 의존성 관리 도구로 pip·Conda도 커스텀 클래스나 서드파티 라이브러리에 그대로 쓸 수 있어요.
다음 단계
- 전체 API를 깊이 있게 보려면 RDD 프로그래밍 가이드와 SQL 프로그래밍 가이드부터 시작하세요.
- 클러스터에서 애플리케이션을 운영하려면 배포 개요(cluster-overview)를 보면 됩니다.
- Spark의
examples디렉터리에 Python/Scala/Java/R 샘플이 여럿 있으니 바로 돌려볼 수 있어요 (./bin/spark-submit examples/src/main/python/pi.py,./bin/run-example SparkPi등).