전체 그래프
Spark

RDD Actions

data-engineeringsparkrddaction

상위: Spark

요약

RDD Action은 실제 연산을 실행하고 결과를 반환하는 연산입니다. collect, count, reduce 등이 있으며, Action이 호출되는 시점에 모든 Transformation이 실행됩니다.

Action 연산 개요

연산설명예제결과
collect()모든 데이터를 리스트로 반환rdd.collect()[1, 2, 3, 4]
count()전체 요소 개수 반환rdd.count()4
reduce()전체 요소를 하나로 결합rdd.reduce(lambda a, b: a + b)10
sum()요소의 합 반환rdd.sum()10
mean()평균 값 반환rdd.mean()2.5

COLLECT

모든 데이터를 Driver로 가져와 리스트로 반환

RDD Actions

x = sc.parallelize([1, 2, 3], 2)
y = x.collect()

print(x.glom().collect())   
print(y)                    

출력:

x: [[1], [2, 3]]
y: [1, 2, 3]

주의: 대용량 데이터에서는 메모리 부족 발생 가능

COUNTBYKEY

각 키별 개수를 세어 딕셔너리로 반환

RDD Actions

x = sc.parallelize([('J', 'James'), ('F', 'Fred'), ('A', 'Anna'), ('J', 'John')])
y = x.countByKey()
print(y) 

출력:

x: [('J', 'James'), ('F', 'Fred'), ('A', 'Anna'), ('J', 'John')]
y: {'J': 2, 'F': 1, 'A': 1}

REDUCE

모든 요소를 하나로 결합

RDD Actions

RDD Actions

RDD Actions

RDD Actions

x = sc.parallelize([1, 2, 3, 4])
y = x.reduce(lambda a, b: a + b)

print(x.collect())
print(y)           

출력:

x: [1, 2, 3, 4]
y: 10

동작 방식:

  1. 파티션 내에서 먼저 reduce 수행 (로컬 reduce)
  2. 각 파티션의 결과를 다시 reduce (글로벌 reduce)

SUM

모든 요소의 합 반환

RDD Actions

x = sc.parallelize([2, 4, 1])
y = x.sum()

print(x.collect())
print(y)           

출력:

x: [2, 4, 1]
y: 7

내부 구현:

self.mapPartitions(lambda x: [sum(x)]).fold(0, operator.add)

MAX

최댓값 반환

RDD Actions

x = sc.parallelize([2, 4, 1])
y = x.max()

print(x.collect())  
print(y)       

출력:

x: [2, 4, 1]
y: 4

내부 구현:

reduce(lambda a, b: max(a, b, key=key))

MEAN

평균값 반환

RDD Actions

x = sc.parallelize([2, 4, 1])
y = x.mean()

print(x.collect())
print(y)       

출력:

x: [2, 4, 1]
y: 2.333333

내부 구현:

self.stats().mean()

STDEV

표준편차 반환

RDD Actions

x = sc.parallelize([2, 4, 1])
y = x.stdev()

print(x.collect()) 
print(y)         

출력:

x: [2, 4, 1]
y: 1.2472191

TAKE

처음 n개 요소 반환

x = sc.parallelize([1, 2, 3, 4, 5])
y = x.take(3)

print(y)

출력:

[1, 2, 3]

FIRST

첫 번째 요소 반환

x = sc.parallelize([1, 2, 3, 4, 5])
y = x.first()

print(y)

출력:

1

FOREACH

각 요소에 함수를 적용 (반환값 없음)

x = sc.parallelize([1, 2, 3])
x.foreach(lambda x: print(x * 2))

출력:

2
4
6

주의: 출력 순서는 보장되지 않음 (병렬 실행)

RDD 데이터 저장하기

saveAsTextFile

텍스트 형식으로 저장

rdd = sc.parallelize([1, 2, 3, 4])
rdd.saveAsTextFile("/path/to/output")

특징:

  • 텍스트 형식으로 저장
  • 사람이 읽기 쉬움
  • 디버깅/분석에 유용
  • 파티션별로 파일 생성 (part-00000, part-00001, ...)

saveAsObjectFile

객체 직렬화 형식으로 저장

rdd = sc.parallelize([{"name": "Alice", "age": 25}, {"name": "Bob", "age": 30}])
rdd.saveAsObjectFile("/path/to/output")

특징:

  • 자바 직렬화 객체로 저장
  • 재사용/공유에 효율적
  • 역직렬화 가능
  • Python 객체는 Pickle 사용

Action 사용 시 주의사항

1. 메모리 관리

  • collect()는 모든 데이터를 Driver 메모리로 가져옴
  • 대용량 데이터에서는 take(n), first() 사용
  • 또는 파일로 저장 후 일부만 읽기

2. 성능 고려

  • Action은 전체 DAG를 실행하므로 비용이 큼
  • 불필요한 Action 호출 최소화
  • 여러 Action이 필요하면 cache() 또는 persist() 사용

3. Shuffle 비용

  • reduce(), countByKey() 등은 Shuffle 발생 가능
  • 네트워크 및 디스크 I/O 비용 고려

추천 패턴

대용량 데이터 확인

# Bad
result = rdd.collect()  # 메모리 부족 위험

# Good
sample = rdd.take(10)   # 일부만 확인

결과 저장

# Bad
result = rdd.collect()
with open("output.txt", "w") as f:
    f.write(str(result))

# Good
rdd.saveAsTextFile("output")  # 분산 저장

반복 사용

# Bad
result1 = rdd.map(...).filter(...).count()
result2 = rdd.map(...).filter(...).sum()  # 같은 연산 반복

# Good
cached = rdd.map(...).filter(...).cache()
result1 = cached.count()
result2 = cached.sum()  # 캐시된 RDD 재사용