전체 그래프
Spark

Data IO

data-engineeringsparkdata-iofile-formats

상위: 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")