PySpark SQL로 DataFrame 다루기

PySpark로 배우는 빅데이터 기초

Upendra Devisetty

Science Analyst, CyVerse

DataFrame API vs SQL 쿼리

  • PySpark에서는 DataFrame API와 SQL 쿼리로 SparkSQL을 사용할 수 있습니다

  • DataFrame API는 데이터 처리를 위한 프로그래밍용 DSL을 제공합니다

  • 변환과 액션을 코드로 구성하기가 더 쉽습니다

  • SQL 쿼리는 간결하고 이해하기 쉬우며 이식성이 높습니다

  • DataFrame 연산은 SQL 쿼리로도 수행할 수 있습니다

PySpark로 배우는 빅데이터 기초

SQL 쿼리 실행

  • SparkSession의 sql() 메서드는 SQL 쿼리를 실행합니다

  • sql() 메서드는 SQL 문을 인자로 받아 결과를 DataFrame으로 반환합니다

df.createOrReplaceTempView("table1")
df2 = spark.sql("SELECT field1, field2 FROM table1")
df2.collect()
[Row(f1=1, f2='row1'), Row(f1=2, f2='row2'), Row(f1=3, f2='row3')]
PySpark로 배우는 빅데이터 기초

데이터 추출용 SQL 쿼리

test_df.createOrReplaceTempView("test_table")
query = '''SELECT Product_ID FROM test_table'''
test_product_df = spark.sql(query)
test_product_df.show(5)
+----------+
|Product_ID|
+----------+
| P00069042|
| P00248942|
| P00087842|
| P00085442|
| P00285442|
+----------+
PySpark로 배우는 빅데이터 기초

SQL로 요약 및 그룹화

test_df.createOrReplaceTempView("test_table")
query = '''SELECT Age, max(Purchase) FROM test_table GROUP BY Age'''
spark.sql(query).show(5)
+-----+-------------+
|  Age|max(Purchase)|
+-----+-------------+
|18-25|        23958|
|26-35|        23961|
| 0-17|        23955|
|46-50|        23960|
|51-55|        23960|
+-----+-------------+
only showing top 5 rows
PySpark로 배우는 빅데이터 기초

SQL로 열 필터링

test_df.createOrReplaceTempView("test_table")
query = '''SELECT Age, Purchase, Gender FROM test_table WHERE Purchase > 20000 AND Gender == "F"'''
spark.sql(query).show(5)
+-----+--------+------+
|  Age|Purchase|Gender|
+-----+--------+------+
|36-45|   23792|     F|
|26-35|   21002|     F|
|26-35|   23595|     F|
|26-35|   23341|     F|
|46-50|   20771|     F|
+-----+--------+------+
only showing top 5 rows
PySpark로 배우는 빅데이터 기초

실습해 봅시다!

PySpark로 배우는 빅데이터 기초

Preparing Video For Download...