상위: Spark
요약
Spark는 wget으로 다운로드하여 설치하고 환경변수를 설정합니다. PySpark는 pip로 간단히 설치 가능하며, spark-shell 또는 pyspark 명령으로 실행을 확인할 수 있습니다.
Spark 설치
wget https://archive.apache.org/dist/spark/spark-3.5.4/spark-3.5.4-bin-hadoop3.tgz
tar -xvzf spark-3.5.4-bin-hadoop3.tgz
sudo mv spark-3.5.4-bin-hadoop3 spark
Spark 환경변수 등록
vi ~/.bashrc
# 아래 두 줄 추가
export SPARK_HOME=/usr/local/spark
export PATH=$SPARK_HOME/bin:$SPARK_HOME/sbin:$PATH
source ~/.bashrc
echo $SPARK_HOME # 확인: /usr/local/spark
Spark 실행 확인
spark-shell
# scala 로고와 함께 version 3.5.4 출력되면 성공
PySpark 설치
pip install pyspark==3.5.4
pyspark # 실행 확인
PySpark 코드 내부 작동 흐름

Python에서 작성한 코드는 다음과 같은 과정을 거쳐 실행됩니다:
- Python 코드 작성 (PySpark API 사용)
- Py4J를 통해 JVM의 Spark API 호출
- Spark Core에서 실행 계획 수립
- Executor에서 실제 작업 수행
- 결과를 Python으로 반환
기본 예제
1. 테스트 파일 생성
echo -e "Hello Spark\nApache Spark is powerful\nBig Data Processing" > test.txt
2. 파일 읽기
rdd = sc.textFile("test.txt")
rdd.foreach(print)
출력:
Hello Spark
Apache Spark is powerful
Big Data Processing
3. 현재 파티션 개수 확인
rdd.getNumPartitions()
# 출력: 2
4. 파티션 재조정
Spark에서 데이터를 하나의 파티션에서 처리:
repartition(1): 파티션 수를 1개로 줄여서 병렬성을 제한하고 데이터를 한곳으로 모음
rdd_repartitioned = rdd.repartition(1)
print("변경된 파티션 수:", rdd_repartitioned.getNumPartitions())
# 출력: 변경된 파티션 수: 1
데이터 변환 기본 예제
1. 데이터 변환 후 결과 확인 (collect)
collect(): RDD의 모든 데이터를 배열 형태로 반환하는 Action 연산- 큰 데이터셋에서는 사용 지양
rdd = sc.textFile("file:///usr/local/spark/data/test.txt")
result = rdd.collect()
print(result)
# 출력: ['Hello Spark', 'Apache Spark is powerful', 'Big Data Processing']
2. 줄 개수 세기 (count)
count(): RDD 전체 줄 개수를 세는 기본적인 Action 연산- 모든 요소의 개수를 반환하며 즉시 평가됨
line_count = rdd.count()
print("줄 개수:", line_count)
# 출력: 줄 개수: 3
3. 특정 단어가 포함된 줄 필터링 (filter)
filter(): 특정 키워드를 포함한 줄만 필터링- Transformation 연산, 조건을 만족하는 요소만 포함하는 새로운 RDD 생성
filtered_rdd = rdd.filter(lambda line: "Spark" in line)
filtered_rdd.foreach(print)
# 출력:
# Hello Spark
# Apache Spark is powerful
4. 모든 단어를 소문자로 변환 (map)
map(): RDD의 각 요소를 변환하는 함수 적용 가능- Transformation 연산, 새로운 RDD 생성
lower_rdd = rdd.map(lambda line: line.lower())
lower_rdd.foreach(print)
# 출력:
# hello spark
# apache spark is powerful
# big data processing
주의사항
메모리 관리
collect()는 모든 데이터를 Driver로 가져오므로 대용량 데이터에서는 메모리 부족 발생 가능- 대신
take(n),first(), 또는 파일로 저장하는 방식 사용
파티션 설정
- 기본 파티션 수는 입력 파일 크기와 블록 크기에 따라 자동 결정
- 너무 적으면 병렬성 저하, 너무 많으면 오버헤드 증가
- 일반적으로
Executor 수 × Core 수 × 2~4정도가 적절
로컬 vs 클러스터 모드
# 로컬 모드 (단일 머신)
spark = SparkSession.builder.master("local[*]").getOrCreate()
# 클러스터 모드 (YARN)
spark = SparkSession.builder.master("yarn").getOrCreate()