Spark DataFrame 더 알아보기

PySpark 입문

Benjamin Schmidt

Data Engineer

여러 데이터 소스에서 DataFrame 생성

  • CSV 파일: 구조화된 구분자 데이터에 일반적
  • JSON 파일: 반정형, 계층형 포맷
  • Parquet 파일: 저장·쿼리에 최적화, 데이터 엔지니어링에 자주 사용
  • 예:
    spark.read.csv("path/to/file.csv")
    
  • 예:
    spark.read.json("path/to/file.json")
    
  • 예:
    spark.read.parquet("path/to/file.parquet")
    
1 https://spark.apache.org/docs/latest/api/python/reference/pyspark.pandas/api/pyspark.pandas.read_csv
PySpark 입문

스키마 추론과 수동 정의

  • inferSchema=True로 데이터에서 스키마 추론 가능

  • 고정 구조엔 수동 스키마 지정 권장 — 더 정확한 제어 가능

대규모 스키마

PySpark 입문

PySpark DataFrame의 DataType

  • IntegerType: 정수
    • 예: 1, 3478, -1890456
  • LongType: 큰 정수
    • 예: 8바이트 부호 있는 수, 922334775806
  • FloatType, DoubleType: 소수 표현 부동소수점
    • 예: 3.14159
  • StringType: 텍스트/문자열 데이터
    • 예: "This is an example of a string."
  • ...
PySpark 입문

PySpark DataFrame DataType 문법

# 필요한 타입 클래스 임포트
from pyspark.sql.types import (StructType,
                            StructField, IntegerType,
                            StringType, ArrayType)

# 스키마 구성
schema = StructType([
    StructField("id", IntegerType(), True),
    StructField("name", StringType(), True),
    StructField("scores", ArrayType(IntegerType()), True)
])

# 스키마 적용
df = spark.createDataFrame(data, schema=schema)
PySpark 입문

DataFrame 작업 - 선택과 필터링

  • 특정 컬럼 선택: .select()
  • 조건으로 행 필터링: .filter() 또는 .where()
  • 여러 컬럼으로 정렬: .sort()
# name, age 컬럼만 선택하여 표시
df.select("name", "age").show()
# age > 30 필터링
df.filter(df["age"] > 30).show()
# 특정 값과 일치하도록 Where 사용
df.where(df["age"] == 30).show()
# age로 정렬
df.sort("age", ascending=False).show()
PySpark 입문

정렬과 결측값 제거

  • .sort() 또는 .orderBy()로 정렬
  • na.drop()으로 null 포함 행 제거
# age 컬럼으로 정렬
df.sort("age", ascending=False).show()

# 결측값 제거
df.na.drop().show()

PySpark 입문

치트시트

  • spark.read_json(): JSON에서 로드
  • spark.read.schema(): 스키마 명시적 지정
  • .na.drop(): 결측값 있는 행 삭제
  • .select(), .filter(), .sort(), .orderBy(): 기본 데이터 조작 함수
PySpark 입문

Let's practice!

PySpark 입문

Preparing Video For Download...