상위: Spark
요약
Spark는 다양한 파일 형식(CSV, JSON, Parquet, Text 등)과 데이터 소스(S3, HDFS, Cassandra 등)를 지원합니다. RDD와 DataFrame 모두 파일 읽기/쓰기를 제공하며, 파티션 단위로 저장됩니다.
RDD 데이터 읽기
1. parallelize (메모리 데이터)
# Python
wordsRDD = sc.parallelize(["fish", "cats", "dogs"])
// Scala
val wordsRDD = sc.parallelize(List("fish", "cats", "dogs"))
// Java
JavaRDD<String> wordsRDD = sc.parallelize(Arrays.asList("fish", "cats", "dogs"))
특징:
- 클러스터로 데이터를 보낼 때 사용
- 비효율적일 수 있어 소규모 분석에 적합
- 테스트 및 프로토타이핑에 유용
2. textFile (외부 파일)
# Python
linesRDD = sc.textFile("/path/to/README.md")
// Scala
val linesRDD = sc.textFile("/path/to/README.md")
// Java
JavaRDD<String> linesRDD = sc.textFile("/path/to/README.md")
지원 소스:
- 로컬 파일 시스템
- HDFS (hdfs://...)
- Amazon S3 (s3://...)
- HBase, Cassandra
- 기타 Hadoop InputFormat
RDD 데이터 저장
saveAsTextFile
rdd = sc.parallelize([1, 2, 3, 4])
rdd.saveAsTextFile("/path/to/output")
특징:
- 텍스트 형식으로 저장
- 사람이 읽기 쉬움
- 디버깅/분석에 유용
- 파티션별로 파일 생성 (
part-00000,part-00001, ...)
출력 예시:
/path/to/output/
├── _SUCCESS
├── part-00000
├── part-00001
└── part-00002
saveAsObjectFile
rdd = sc.parallelize([{"name": "Alice", "age": 25}, {"name": "Bob", "age": 30}])
rdd.saveAsObjectFile("/path/to/output")
특징:
- 자바 직렬화 객체로 저장
- 재사용/공유에 효율적
- 역직렬화 가능
- Python 객체는 Pickle 사용
DataFrame 데이터 읽기
CSV
# 기본 읽기
df = spark.read.csv("data.csv")
# 옵션 포함
df = spark.read \
.option("header", "true") \
.option("inferSchema", "true") \
.option("delimiter", ",") \
.csv("data.csv")
# 단축 형식
df = spark.read.csv("data.csv", header=True, inferSchema=True)
JSON
# 기본 읽기
df = spark.read.json("data.json")
# 여러 파일
df = spark.read.json(["file1.json", "file2.json"])
# 멀티라인 JSON
df = spark.read.option("multiLine", "true").json("data.json")
Parquet
# Parquet 읽기 (컬럼 기반 저장)
df = spark.read.parquet("data.parquet")
# 여러 파일
df = spark.read.parquet("data.parquet", "data2.parquet")
Parquet 장점:
- 컬럼 기반 저장으로 압축률 높음
- 필요한 컬럼만 읽기 가능 (효율적)
- 스키마 자동 포함
- Big Data 표준 포맷
Text
# 텍스트 파일
df = spark.read.text("data.txt")
기타 형식
# ORC
df = spark.read.orc("data.orc")
# Avro
df = spark.read.format("avro").load("data.avro")
DataFrame 데이터 쓰기
CSV
df.write.csv("/path/to/output")
# 옵션 포함
df.write \
.option("header", "true") \
.option("delimiter", "|") \
.csv("/path/to/output")
JSON
df.write.json("/path/to/output")
# 모드 지정
df.write.mode("overwrite").json("/path/to/output")
Parquet
df.write.parquet("/path/to/output")
# 파티셔닝
df.write.partitionBy("year", "month").parquet("/path/to/output")
출력 예시 (파티셔닝):
/path/to/output/
├── year=2024/
│ ├── month=01/
│ │ └── part-00000.parquet
│ └── month=02/
│ └── part-00000.parquet
└── _SUCCESS
저장 모드
# overwrite: 기존 데이터 덮어쓰기
df.write.mode("overwrite").parquet("/path/to/output")
# append: 기존 데이터에 추가
df.write.mode("append").parquet("/path/to/output")
# ignore: 이미 존재하면 무시
df.write.mode("ignore").parquet("/path/to/output")
# error (기본값): 이미 존재하면 에러
df.write.mode("error").parquet("/path/to/output")
외부 데이터 소스
JDBC (데이터베이스)
# 읽기
df = spark.read \
.format("jdbc") \
.option("url", "jdbc:postgresql://localhost/test") \
.option("dbtable", "users") \
.option("user", "username") \
.option("password", "password") \
.load()
# 쓰기
df.write \
.format("jdbc") \
.option("url", "jdbc:postgresql://localhost/test") \
.option("dbtable", "users") \
.option("user", "username") \
.option("password", "password") \
.save()
S3
# S3 읽기
df = spark.read.parquet("s3a://bucket-name/path/to/data.parquet")
# S3 쓰기
df.write.parquet("s3a://bucket-name/path/to/output")
HDFS
# HDFS 읽기
df = spark.read.parquet("hdfs://namenode:8020/path/to/data.parquet")
# HDFS 쓰기
df.write.parquet("hdfs://namenode:8020/path/to/output")
파티션 관리
파티션 수 확인
# RDD
rdd.getNumPartitions()
# DataFrame
df.rdd.getNumPartitions()
파티션 수 조정
# repartition: 파티션 수 변경 (Shuffle 발생)
df_repartitioned = df.repartition(10)
# coalesce: 파티션 수 줄이기 (Shuffle 최소화)
df_coalesced = df.coalesce(5)
파일 크기 제어
# 작은 파일 방지: 파티션 수 줄이기
df.coalesce(1).write.parquet("/path/to/output")
# 큰 파일 분할: 파티션 수 늘리기
df.repartition(100).write.parquet("/path/to/output")
데이터 압축
Parquet 압축
# snappy (기본값)
df.write.option("compression", "snappy").parquet("/path/to/output")
# gzip
df.write.option("compression", "gzip").parquet("/path/to/output")
# none
df.write.option("compression", "none").parquet("/path/to/output")
CSV/JSON 압축
# gzip 압축
df.write.option("compression", "gzip").json("/path/to/output")
스트리밍 I/O
스트리밍 읽기
# 파일 스트리밍
streaming_df = spark.readStream \
.schema(schema) \
.csv("/path/to/streaming/input")
스트리밍 쓰기
query = streaming_df.writeStream \
.outputMode("append") \
.format("parquet") \
.option("path", "/path/to/output") \
.option("checkpointLocation", "/path/to/checkpoint") \
.start()
성능 최적화 팁
1. 적절한 파일 형식 선택
- Parquet: 분석 워크로드, 컬럼 기반 쿼리
- ORC: Hive 통합, 트랜잭션
- CSV: 호환성, 간단한 데이터
- JSON: 반구조화 데이터
2. 파티셔닝 활용
# 시간 기반 파티셔닝
df.write.partitionBy("year", "month", "day").parquet("/path/to/output")
# 필터링 시 파티션 프루닝
df = spark.read.parquet("/path/to/output")
filtered = df.filter(col("year") == 2024) # 2024 파티션만 읽음
3. 스키마 명시
# Bad: inferSchema는 전체 파일 스캔
df = spark.read.csv("data.csv", inferSchema=True)
# Good: 스키마 명시
df = spark.read.csv("data.csv", schema=predefined_schema)
4. 적절한 파티션 수
# 파일 크기 128MB ~ 1GB 권장
# 너무 작은 파일: 메타데이터 오버헤드
# 너무 큰 파일: 메모리 부족, 병렬성 저하
# 적절한 파티션 수 = 전체 데이터 크기 / 목표 파티션 크기
num_partitions = total_size_mb / 256 # 256MB per partition
df.repartition(num_partitions).write.parquet("/path/to/output")