Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
月別売上・分類別順位・累積和を SQL と DataFrame で
한국어 원문으로 표시합니다.
목표
같은 질문을 Spark SQL 과 DataFrame API 로 각각 풀어 같은 물리 계획이 나오는 것을 확인하고, 윈도 함수로 순위·누적합·전월 대비를 계산한다. 마지막으로 정확한 고유 수와 근사 고유 수를 견준다.
왜 중요한가
Spark 에서 SQL 과 DataFrame 은 두 개의 엔진이 아니다. 둘 다 같은 논리 계획으로 바뀌고, 같은 옵티마이저(Catalyst)를 거쳐, 같은 물리 계획으로 실행된다. 그래서 "SQL 이 더 빠르다" 나 "DataFrame 이 더 빠르다" 는 대개 틀린 질문이다. 팀이 무엇으로 쓰든 계획을 보면 같은지 다른지 바로 안다. 집계 다음에 가장 자주 쓰는 것이 윈도 함수다. 행을 줄이지 않고 옆 행을 볼 수 있어서 순위·누적합·전월 대비 같은 보고서 숫자가 전부 여기서 나온다. 윈도의 틀(파티션·정렬·범위)을 잘못 잡으면 오류 없이 틀린 숫자가 나오므로, 틀을 말로 설명할 수 있어야 한다. 고유 수 세기는 셔플이 큰 연산이다. 몇 퍼센트의 오차를 받아들이면 HyperLogLog++ 로 훨씬 싸게 셀 수 있다. 대시보드에는 근사, 청구서에는 정확 — 어느 쪽을 쓸지는 숫자의 쓰임이 정한다.
단계
- /root/spk/sql/common.py 에 주문·상품을 스키마로 읽어 임시 뷰
orders·products로 등록하는load(spark)를 만들고, /root/spk/sql/sql_monthly.py(앱spk-sql-monthly)에서 SQL 로 월별 매출(칼럼month,revenue)을 /root/spk/sql/out/monthly_sql 에 머리줄 있는 CSV 로 쓰세요. - /root/spk/sql/df_monthly.py(앱
spk-sql-df)에서spark.sql없이 DataFrame API 로 같은 결과를 /root/spk/sql/out/monthly_df 에 쓰세요. - /root/spk/sql/plans.py(앱
spk-sql-plans)에서 두 방식의explain(mode="formatted")출력을 /root/spk/sql/out/plan_sql.txt·/root/spk/sql/out/plan_df.txt 에 저장한 뒤, 두 DataFrame 을 각각collect()로 실제로 돌리세요. - /root/spk/sql/top3.py(앱
spk-sql-top3)에서 분류(category)마다 매출 상위 3 상품을 /root/spk/sql/out/top3 에 칼럼category,product_id,revenue,rank로 쓰세요. 동점이면product_id가 앞선 쪽이 위입니다. - /root/spk/sql/running.py(앱
spk-sql-running)에서 2026년 1월의 채널별 일 매출과 채널 안 누적합을 /root/spk/sql/out/running 에 칼럼channel,day,revenue,running으로 쓰세요. - /root/spk/sql/mom.py(앱
spk-sql-mom)에서 채널별 월 매출과 전월 매출, 증가율(퍼센트, 소수 둘째 자리 반올림)을 /root/spk/sql/out/mom 에 칼럼channel,month,revenue,prev_revenue,growth_pct로 쓰세요. 첫 달의 전월 값은 비워 둡니다. - /root/spk/sql/distinct.py(앱
spk-sql-distinct)에서 결제 완료 주문을 한 고객 수를 정확히(countDistinct), 그리고approx_count_distinct(rsd=0.05)로 세어 /root/spk/sql/out/distinct.json 에{"exact": 정수, "approx": 정수}로 쓰세요. - /root/spk/sql/report.md 에
## SQL 과 DataFrame## 윈도 함수## 근사 집계세 절을 쓰세요. 셋째 절에는 7단계의 두 숫자를 넣으세요.
참고
- 원본:
/data/shop/orders.csv(order_id, customer_id, product_id, qty, order_ts, status, channel),/data/shop/products.csv(product_id, category, price). 매출 = 결제 완료(status='paid') 주문의qty × price입니다. - 스크립트를
common.py와 같은 폴더(/root/spk/sql)에 두면from common import load가 됩니다(실행한 스크립트의 폴더가 파이썬 모듈 경로에 들어갑니다). - 결과 CSV 는 작으니
coalesce(1)로 파일 하나에 모으면 눈으로 보기 쉽습니다(채점기는 파일이 여럿이어도 읽습니다). explain은 표준출력으로 찍습니다.contextlib.redirect_stdout으로 받아 파일에 쓰세요.- 흔한 실수:
rank를 써서 동점에 같은 순위가 둘 나오는 것, 누적합의 창을 기본값(정렬이 있으면 RANGE … CURRENT ROW)으로 두어 같은 날이 겹치는 것, 취소·환불 주문을 매출에 넣는 것. - 공식 문서: Spark SQL Guide · Window Functions · EXPLAIN · Built-in Functions
임시 뷰와 SQL 로 월별 매출
/root/spk/sql/common.py 에 주문·상품을 스키마로 읽어 임시 뷰 orders·products 로 등록하는 load(spark) 를 만들고, /root/spk/sql/sql_monthly.py 를 앱 이름 spk-sql-monthly 로 만들어 spark.sql 로 결제 완료 주문의 월별 매출(칼럼 month 는 yyyy-MM, revenue 는 qty×price 의 합)을 /root/spk/sql/out/monthly_sql 에 머리줄 있는 CSV 로 쓰세요.
임시 뷰는 이 SparkSession 안에서만 보이는 이름입니다. 뷰를 등록해도 아무것도 읽지 않습니다 — 계획에 이름을 붙일 뿐입니다. date_format(order_ts, 'yyyy-MM') 으로 달을 만들고 상품과 조인해 가격을 곱하세요.
같은 질문을 DataFrame API 로
/root/spk/sql/df_monthly.py 를 앱 이름 spk-sql-df 로 만들어 spark.sql 을 쓰지 않고 DataFrame API(where·join·groupBy·agg)로 1단계와 같은 결과를 /root/spk/sql/out/monthly_df 에 머리줄 있는 CSV(month,revenue)로 쓰세요.
F.date_format("order_ts", "yyyy-MM").alias("month") 로 묶고 F.sum(F.col("qty") * F.col("price")) 로 더합니다. 조인 키 이름이 같으면 join(p, "product_id") 처럼 문자열로 주어 칼럼이 하나만 남게 하세요.
두 방식의 물리 계획이 같은지 보기
/root/spk/sql/plans.py 를 앱 이름 spk-sql-plans 로 만들어 1·2단계의 두 질의를 각각 만들고 explain(mode="formatted") 출력을 /root/spk/sql/out/plan_sql.txt 와 /root/spk/sql/out/plan_df.txt 에 저장한 다음, 두 DataFrame 을 각각 collect() 로 실제로 돌리세요.
formatted 모드는 위에 연산자 트리, 아래에 번호 붙은 연산자 설명을 찍습니다. 칼럼 번호(#12 같은 것)는 질의마다 달라지지만 연산자의 이름과 순서는 같아야 합니다. 실제로 돌리면 이벤트 로그의 SQL 실행 기록에 똑같은 계획 문서가 남습니다 — 채점기는 여러분의 파일이 그 기록과 같은지, 두 계획의 연산자가 같은지를 봅니다.
분류별 상위 3 상품 — 윈도 순위
/root/spk/sql/top3.py 를 앱 이름 spk-sql-top3 로 만들어 분류(category)·상품별 결제 완료 매출을 구하고, 분류마다 매출 내림차순(동점이면 product_id 오름차순)으로 1–3위를 /root/spk/sql/out/top3 에 머리줄 있는 CSV(category,product_id,revenue,rank)로 쓰세요.
윈도는 Window.partitionBy("category").orderBy(...) 입니다. rank() 는 동점에 같은 번호를 주고 다음 번호를 건너뛰므로, 정확히 3개를 원하면 row_number() 에 동점 규칙을 정렬로 넣어야 합니다.
채널별 누적합 — 창의 범위 정하기
/root/spk/sql/running.py 를 앱 이름 spk-sql-running 으로 만들어 2026년 1월 결제 완료 주문의 채널별 일 매출(day 는 날짜)과, 채널 안에서 날짜 순으로 더한 누적합 running 을 /root/spk/sql/out/running 에 머리줄 있는 CSV(channel,day,revenue,running)로 쓰세요.
누적합은 Window.partitionBy("channel").orderBy("day").rowsBetween(Window.unboundedPreceding, Window.currentRow) 위의 sum 입니다. 날짜마다 한 줄로 먼저 모은 뒤 창을 씌우세요 — 모으기 전에 씌우면 같은 날의 주문마다 누적합이 달라집니다.
전월 대비 — lag 로 옆 행 보기
/root/spk/sql/mom.py 를 앱 이름 spk-sql-mom 으로 만들어 채널별 월 매출과 lag 로 가져온 전월 매출 prev_revenue, 증가율 growth_pct(=(이번 달-전월)/전월×100, 소수 둘째 자리 반올림)를 /root/spk/sql/out/mom 에 머리줄 있는 CSV(channel,month,revenue,prev_revenue,growth_pct)로 쓰세요. 채널의 첫 달은 전월·증가율이 비어 있어야 합니다.
F.lag("revenue").over(Window.partitionBy("channel").orderBy("month")) 가 같은 채널의 바로 앞 달 값을 가져옵니다. 첫 달은 앞 행이 없어 null 이고, null 이 섞인 나눗셈도 null 이라 따로 처리할 필요가 없습니다.
정확한 고유 수와 근사 고유 수
/root/spk/sql/distinct.py 를 앱 이름 spk-sql-distinct 로 만들어 결제 완료 주문을 한 고객 수를 countDistinct 와 approx_count_distinct(rsd=0.05) 로 한 번에 세어 /root/spk/sql/out/distinct.json 에 {"exact": 정수, "approx": 정수} 로 쓰세요.
근사 쪽은 HyperLogLog++ 스케치를 합칩니다. 고유 값을 셔플로 모으지 않아도 되니 커질수록 싸집니다. rsd 는 상대 표준 오차의 목표치라서, 결과가 정확한 값과 몇 퍼센트 다른지 직접 계산해 보세요. 채점기는 이벤트 로그의 계획에 근사 함수가 실제로 들어 있는지도 봅니다.
같은 계획·윈도의 틀·근사의 오차를 남기기
/root/spk/sql/report.md 에 ## SQL 과 DataFrame ## 윈도 함수 ## 근사 집계 세 절을 쓰세요. 셋째 절에는 7단계의 정확한 값과 근사 값을 숫자로 넣으세요.
첫 절에는 두 계획에서 같았던 연산자 이름을 몇 개 적고, 둘째 절에는 순위·누적합·전월 대비에서 창을 어떻게 잡았는지 한 줄씩 적으세요. 셋째 절에는 두 숫자와 그 차이의 퍼센트를 적으면 됩니다.