전체 그래프
Spark

Installation

data-engineeringsparksetuppyspark

상위: 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 코드 내부 작동 흐름

Installation

Python에서 작성한 코드는 다음과 같은 과정을 거쳐 실행됩니다:

  1. Python 코드 작성 (PySpark API 사용)
  2. Py4J를 통해 JVM의 Spark API 호출
  3. Spark Core에서 실행 계획 수립
  4. Executor에서 실제 작업 수행
  5. 결과를 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()