PySpark 입문
Benjamin Schmidt
Data Engineer
spark.sql()로 사용할 수 있습니다# SQL 집계 쿼리
spark.sql("""
SELECT Department, SUM(Salary) AS Total_Salary, AVG(Salary) AS Average_Salary
FROM employees
GROUP BY Department
""").show()
# 3000 초과 급여 필터링 filtered_df = df.filter(df.Salary > 3000) # 필터링된 DataFrame을 뷰로 등록 filtered_df.createOrReplaceTempView("filtered_employees")# 필터링된 뷰에 SQL로 집계 수행 spark.sql(""" SELECT Department, COUNT(*) AS Employee_Count FROM filtered_employees GROUP BY Department """).show()
# 타입 캐스팅 예시 data = [("HR", "3000"), ("IT", "4000"), ("Finance", "3500")] columns = ["Department", "Salary"] df = spark.createDataFrame(data, schema=columns)# Salary 컬럼을 정수로 변환 df = df.withColumn("Salary", df["Salary"].cast("int")) # 집계 수행 df.groupBy("Department").sum("Salary").show()
# RDD로 집계 예시 rdd = df.rdd.map(lambda row: (row["Department"], row["Salary"]))rdd_aggregated = rdd.reduceByKey(lambda x, y: x + y)print(rdd_aggregated.collect())
groupBy() 같은 연산 최소화explain()으로 실행 계획 확인 및 최적화SUM(), AVERAGE() 등의 함수PySpark 입문