PySpark 에서 데이터를 읽고 쓰는 기본 형태를 모은다.
df = (spark.read
.format("parquet")
.load("s3a://bucket/path/"))
포맷별 축약형도 있다.
spark.read.parquet("s3a://bucket/path/")
spark.read.json("s3a://bucket/events/")
spark.read.option("header", True).option("inferSchema", True).csv("s3a://bucket/raw.csv")
inferSchema 는 파일을 한 번 더 읽는다. 큰 데이터에서는 스키마를 직접 준다.
from pyspark.sql.types import StructType, StructField, StringType, LongType, TimestampType
schema = StructType([
StructField("id", LongType(), False),
StructField("name", StringType(), True),
StructField("created_at", TimestampType(), True),
])
df = spark.read.schema(schema).json("s3a://bucket/events/")
(df.write
.format("parquet")
.mode("overwrite")
.partitionBy("dt")
.save("s3a://bucket/output/"))
| mode | 동작 |
|---|---|
errorifexists |
기본값. 경로가 있으면 실패 |
append |
이어서 쓴다 |
overwrite |
경로를 지우고 다시 쓴다 |
ignore |
경로가 있으면 아무것도 하지 않는다 |
overwrite는 파티션 단위가 아니라 경로 전체를 지운다. 특정 파티션만 바꾸려면 아래 설정을 켠다.
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
df = (spark.read
.format("jdbc")
.option("url", "jdbc:postgresql://db:5432/app")
.option("dbtable", "public.orders")
.option("user", "reader")
.option("password", "<비밀번호>")
.option("driver", "org.postgresql.Driver")
.load())
병렬로 읽으려면 분할 기준을 준다. 안 주면 executor 하나가 전부 가져온다.
.option("partitionColumn", "id")
.option("lowerBound", 1)
.option("upperBound", 10000000)
.option("numPartitions", 8)
df.write.format("delta").mode("overwrite").save("/mnt/delta/orders")
(spark.read.format("delta")
.option("versionAsOf", 3)
.load("/mnt/delta/orders"))
df.coalesce(1).write.mode("overwrite").parquet("s3a://bucket/small/")
df.repartition(16).write.mode("overwrite").parquet("s3a://bucket/big/")
coalesce 는 셔플 없이 파티션을 줄이고, repartition 은 셔플을 한다. 크기가 큰 데이터에 coalesce(1) 을 쓰면 executor 하나에 다 몰린다.