Spark SQL 시작하기

Spark SQL 시작하기 (Getting Started)

Spark SQL을 쓰는 모든 흐름의 출발점이 되는 진입점은 SparkSession이에요. DataFrame을 만들고, 그 위에서 SQL을 실행하고, 임시 뷰를 등록하는 일련의 패턴을 여기서 익히면 그다음 데이터 소스나 성능 튜닝으로 자연스럽게 이어집니다. 구조화된 데이터를 다루는 Spark의 가장 기본적인 워크플로우를 함께 해봐요.

출처: Apache Spark 공식 문서 – SQL Getting Started

시작점: SparkSession

Spark의 모든 기능으로 들어가는 진입점은 SparkSession 클래스예요. 기본 SparkSession은 SparkSession.builder로 만들면 됩니다.

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("...").getOrCreate()

Spark 2.0부터 SparkSession은 Hive 기능에 대한 내장 지원을 제공해요. HiveQL로 쿼리를 작성하거나, Hive UDF에 접근하거나, Hive 테이블의 데이터를 읽는 기능을 말해요. 이 기능을 쓰기 위해 기존 Hive 설정이 있을 필요는 없습니다.

DataFrame 만들기 (Creating DataFrames)

SparkSession이 있으면 애플리케이션은 기존 RDD에서, Hive 테이블에서, 또는 Spark 데이터 소스에서 DataFrame을 만들 수 있어요. 예를 들어 JSON 파일의 내용을 기반으로 DataFrame을 만듭니다.

df = spark.read.json("examples/src/main/resources/people.json")

# DataFrame 스키마 출력
df.printSchema()

비타입 Dataset 연산 (aka DataFrame 연산)

DataFrame은 Python·Scala·Java·R에서 구조화된 데이터를 다루는 도메인 특화 언어(DSL)를 제공해요. Spark 2.0에서 DataFrame은 Scala/Java API에서 Row들의 Dataset일 뿐이라, 강타입 Scala/Java Dataset이 주는 "typed transformation"과 대비해 "untyped transformation"이라고도 불러요.

Python에서는 DataFrame의 컬럼에 df.age처럼 속성으로, 또는 df['age']처럼 인덱싱으로 접근할 수 있어요. 전자가 인터랙티브 탐색에는 편하지만, DataFrame 클래스의 속성과 겹치는 컬럼명에서 깨지지 않는 후자 방식(df['age'])을 공식 문서는 더 권장합니다.

# 컬럼 선택하고 조건 필터링
df.select("name", "age").show()
df.filter(df.age > 21).show()

# 그룹별 집계
df.groupBy("age").count().show()

간단한 컬럼 참조와 표현식 외에도 DataFrame은 문자열 조작, 날짜 연산, 일반 수학 연산 등을 포함한 풍부한 함수 라이브러리를 갖고 있어요. 전체 목록은 DataFrame Function Reference에 있습니다.

프로그래밍 방식으로 SQL 쿼리 실행하기

SparkSession의 sql 함수는 애플리케이션에서 SQL 쿼리를 프로그래밍 방식으로 실행하고 결과를 DataFrame으로 반환하게 해줘요.

# 임시 뷰 등록
df.createOrReplaceTempView("people")

# SQL 쿼리 실행
sqlDF = spark.sql("SELECT name FROM people WHERE age BETWEEN 13 AND 19")
sqlDF.show()

전역 임시 뷰 (Global Temporary View)

Spark SQL의 임시 뷰는 세션 스코프예요. 뷰를 만든 세션이 종료되면 사라집니다. 모든 세션이 공유하고 Spark 애플리케이션이 끝날 때까지 유지되는 임시 뷰를 만들고 싶다면 전역 임시 뷰를 쓰면 돼요. 전역 임시 뷰는 시스템 보존 데이터베이스 global_temp에 연결되며, SELECT * FROM global_temp.view1처럼 정규화된 이름으로 참조해야 합니다.

df.createGlobalTempView("people")

spark.sql("SELECT * FROM global_temp.people").show()

SQL로는 이렇게 됩니다.

CREATE GLOBAL TEMPORARY VIEW temp_view AS SELECT a + 1, b * 2 FROM tbl

SELECT * FROM global_temp.temp_view

Dataset 만들기 (Creating Datasets)

Dataset은 RDD와 비슷하지만, Java 직렬화나 Kryo 대신 전문화된 Encoder로 객체를 직렬화해요. Encoder와 표준 직렬화 모두 객체를 바이트로 바꾸는 책임을 지지만, Encoder는 동적으로 코드 생성되며 필터링·정렬·해싱 같은 많은 연산을 바이트를 다시 객체로 역직렬화하지 않고 수행하게 해주는 포맷을 사용합니다. (Dataset 생성은 Scala/Java API에서 제공돼요.)

RDD와 상호 운용 (Interoperating with RDDs)

Spark SQL은 기존 RDD를 Dataset으로 변환하는 두 가지 방법을 지원해요. 첫 번째는 RDD가 담은 특정 타입 객체의 스키마를 리플렉션으로 추론하는 방법입니다. 이 방식은 코드가 더 간결하고, Spark 애플리케이션을 작성할 때 스키마를 이미 알고 있을 때 잘 맞아요. 두 번째는 프로그래밍 방식으로 스키마를 구성해 기존 RDD에 적용하는 방법인데, 다소 장황하지만 컬럼과 타입이 런타임에야 확정되는 상황에서 유용합니다.

리플렉션으로 스키마 추론하기

Spark SQL은 Row 객체들의 RDD를 데이터타입을 추론하며 DataFrame으로 변환할 수 있어요. Python에서 Row는 키/값 쌍 리스트를 kwargs로 넘겨 만들고, 이 키들이 테이블의 컬럼명이 되며 JSON 파일 추론과 유사하게 전체 데이터셋 샘플링으로 타입을 추론합니다.

from pyspark.sql import Row

# Row 객체의 RDD 만들기
peopleRDD = spark.sparkContext.parallelize([
    Row(name="Andy", age=30),
    Row(name="Justin", age=19),
])

# 스키마 추론해서 DataFrame으로 변환
peopleDF = spark.createDataFrame(peopleRDD)
peopleDF.printSchema()

Scala에서는 case class가 들어있는 RDD를 자동으로 DataFrame으로 변환해 줘요. case class가 테이블의 스키마를 정의하고, 인자 이름은 리플렉션으로 읽혀 컬럼명이 됩니다. case class는 중첩되거나 Seq, Array 같은 복합 타입도 가질 수 있어요. Java에서는 JavaBean의 BeanInfo가 스키마를 정의합니다 (단, Map 필드를 가진 JavaBean은 현재 지원하지 않아요).

프로그래밍 방식으로 스키마 지정하기

kwargs 사전을 미리 정의할 수 없을 때(예: 레코드 구조가 문자열로 인코딩되어 있거나, 텍스트 데이터셋을 파싱해 사용자마다 다른 필드를 투영해야 할 때) DataFrame은 세 단계로 프로그래밍 방식으로 만들어집니다.

  1. 원래 RDD에서 튜플/리스트의 RDD를 만든다.
  2. 1단계의 RDD 구조와 일치하는 StructType으로 스키마를 만든다.
  3. SparkSession이 제공하는 createDataFrame 메서드로 스키마를 RDD에 적용한다.
from pyspark.sql.types import StructType, StructField, StringType, LongType

# 1. 튜플의 RDD
peopleRDD = spark.sparkContext.parallelize([("Justin", 19), ("Andy", 30)])

# 2. 스키마 정의
schema = StructType([
    StructField("name", StringType(), True),
    StructField("age", LongType(), True),
])

# 3. 스키마 적용
peopleDF = spark.createDataFrame(peopleRDD, schema)
peopleDF.printSchema()

스칼라 함수와 집계 함수

Spark SQL은 다양한 내장 스칼라 함수사용자 정의 스칼라 함수(UDF)를 지원해요. 스칼라 함수는 행 그룹에 대해 값을 반환하는 집계 함수와 달리, 행마다 단일 값을 반환하는 함수예요.

집계 함수는 행 그룹에 대해 단일 값을 반환합니다. 내장 집계 함수count(), count_distinct(), avg(), max(), min() 같은 일반적인 집계를 제공하고, 필요한 경우 사용자 정의 집계 함수(UDAF)도 직접 만들 수 있어요.

더 알아보기