LabHub
배우기 러닝패스 코스

Apache Flink — Running Streams on a Real Engine

Read the plan and predict the vertex count

LabHub 에서 이어서 보기

한국어 원문으로 표시합니다.

목표

같은 GROUP BY 질의의 실행 계획을 EXPLAIN 으로 받아 연산자와 Exchange 의 자리를 읽고, 실제 잡의 정점 수를 체이닝과 전달 방식으로 설명한다. mini-batch 두 단계 집계와 distinct 분할이 계획에 어떻게 나타나는지 확인한다.

왜 중요한가

느린 잡의 원인 — 원천까지 내려가지 못한 필터, 한 키에 몰린 집계, 끊긴 체인 — 은 처리량 그래프에서는 똑같이 보이고 계획에서는 다르게 보인다. 튜닝 설정도 계획이 바뀌어야 먹은 것이다. 이 실습의 채점기는 클러스터에 묻지 않는다. 여러분이 저장한 EXPLAIN 출력 원문REST 응답만 읽고, 정점 수는 JSON 실행 계획의 전달 방식과 교차 대조한다.

단계

  1. flink-up/root/flink/plan/ddl.sql 에 원천 clicks 와 blackhole 싱크 page_stats 를 만들고, /root/flink/plan/explain.sqlEXPLAIN PLAN FOR 출력을 /root/flink/plan/explain.out 에 저장하세요.
  2. explain.out 의 실행 계획을 읽어 /root/flink/plan/shuffle.jsonexchange·above·below·filter_pushed_into_scan 을 적으세요.
  3. EXPLAIN ESTIMATED_COST, PLAN_ADVICE 를 쓴 /root/flink/plan/advice.sql 의 출력을 /root/flink/plan/advice.out 에 저장하세요.
  4. EXPLAIN JSON_EXECUTION_PLAN INSERT INTO page_stats … 를 쓴 /root/flink/plan/json.sql 의 출력을 /root/flink/plan/json.out 에 저장하세요.
  5. 잡 이름 flk-plan-chained 로 같은 INSERT 를 끝까지 돌리는 /root/flink/plan/chained.sql 을 실행하고 그 잡의 /jobs/<jid>/root/flink/plan/chained-job.json 에 저장하세요.
  6. 체이닝을 끄고 잡 이름 flk-plan-unchained 로 돌리는 /root/flink/plan/unchained.sql 을 실행하고 /jobs/<jid>/root/flink/plan/unchained-job.json 에 저장하세요.
  7. mini-batch 세 설정을 준 뒤 같은 SELECT 를 EXPLAIN 하는 /root/flink/plan/twophase.sql 의 출력을 /root/flink/plan/twophase.out 에 저장하세요.
  8. distinct 분할(버킷 64)을 켜고 COUNT(DISTINCT user_id) 를 EXPLAIN 하는 /root/flink/plan/distinct.sql 의 출력을 /root/flink/plan/distinct.out 에 저장하고, 숫자를 모은 /root/flink/plan/report.json 을 쓰세요.

참고

EXPLAIN 의 세 구역을 받는다

flink-up/root/flink/plan/ddl.sqlplan_clicks.csv 를 읽는 clicksCREATE TABLE page_stats (page STRING, n_views BIGINT, revenue BIGINT) WITH ('connector' = 'blackhole') 를 쓰고, /root/flink/plan/explain.sqlEXPLAIN PLAN FOR + 참고의 질의를 써서 sql-client.sh -i ddl.sql -f explain.sql > explain.out 2>&1/root/flink/plan/explain.out 을 만드세요.

EXPLAIN 은 잡을 돌리지 않고 계획만 돌려줍니다. 출력에 Abstract Syntax Tree · Optimized Physical Plan · Optimized Execution Plan 세 구역이 보여야 합니다. 초기화 파일(-i)에는 CREATE TABLE 두 문장을 한 줄에 하나씩 둡니다.

섞는 자리를 찾는다

explain.out 의 실행 계획 구역을 읽어 /root/flink/plan/shuffle.jsonexchange(Exchange 의 distribution 값, 예: hash[…]), above(Exchange 보다 위 연산자 이름 목록), below(아래 연산자 이름 목록, 위에서부터), filter_pushed_into_scan(원천 줄에 filter=[…] 가 붙었는가, 참/거짓)을 적으세요.

연산자 이름은 괄호 앞 낱말입니다(GroupAggregate, Calc 등). 트리는 위가 싱크 쪽입니다. WHERE 가 원천으로 내려갔다면 TableSourceScan 줄 안에 filter=[…] 가 보입니다. 채점기는 여러분의 explain.out 과 대조합니다.

비용 추정과 조언을 붙인다

EXPLAIN ESTIMATED_COST, PLAN_ADVICE + 같은 질의를 쓴 /root/flink/plan/advice.sql 을 돌려 /root/flink/plan/advice.out 에 저장하세요. 비용(cumulative cost)과 advice[1]: [ADVICE] 줄이 보여야 합니다.

PLAN_ADVICE 를 붙이면 물리 계획 구역의 제목이 'With Advice' 로 바뀌고 끝에 조언 줄이 붙습니다. 무엇을 켜 보라고 하는지 읽어 두세요 — 7단계에서 씁니다. 추정 비용은 통계 없는 원천에 대한 가정값입니다.

JSON 실행 계획에서 전달 방식을 본다

EXPLAIN JSON_EXECUTION_PLAN INSERT INTO page_stats + 같은 질의를 쓴 /root/flink/plan/json.sql 을 돌려 /root/flink/plan/json.out 에 저장하세요. 싱크(Writer)까지 포함한 노드와 ship_strategy 가 보여야 합니다.

SELECT 가 아니라 INSERT 를 EXPLAIN 해야 실제로 돌릴 잡과 같은 그래프가 나옵니다. == Physical Execution Plan == 아래의 nodes 배열이 연산자이고, predecessors 의 ship_strategy 가 앞 연산자에서 데이터를 받는 방식입니다. HASH 가 몇 번 나오는지 세어 두세요.

체이닝을 켠 잡의 정점을 센다

SET 'pipeline.name' = 'flk-plan-chained'; · SET 'table.dml-sync' = 'true'; · INSERT INTO page_stats + 같은 질의를 쓴 /root/flink/plan/chained.sql 을 돌린 뒤, 그 잡의 /jobs/<jid>/root/flink/plan/chained-job.json 에 저장하세요. 정점 수는 json.out 의 FORWARD 가 아닌 전달 수 + 1 이어야 합니다.

table.dml-sync 를 켜면 sql-client 가 INSERT 잡이 끝날 때까지 기다려, 끝난(FINISHED) 잡의 응답을 받을 수 있습니다. 정점 이름에 -> 로 이어진 연산자들이 한 태스크로 묶인 체인입니다.

체이닝을 끄고 다시 센다

chained.sql 에 SET 'pipeline.operator-chaining.enabled' = 'false'; 를 더하고 잡 이름을 flk-plan-unchained 로 바꾼 /root/flink/plan/unchained.sql 을 돌린 뒤 /jobs/<jid>/root/flink/plan/unchained-job.json 에 저장하세요. 정점 수가 json.out 의 연산자(노드) 수와 같아야 합니다.

체이닝을 끄면 FORWARD 로 이어진 연산자들도 각자 정점(태스크)이 됩니다. 결과는 같고 태스크 사이 넘김만 늘어납니다. 이름을 바꾸지 않으면 잡 목록에서 앞 단계의 잡과 헷갈립니다.

mini-batch 로 두 단계 집계를 만든다

table.exec.mini-batch.enabled = true, table.exec.mini-batch.allow-latency = 1 s, table.exec.mini-batch.size = 1000 을 SET 한 뒤 EXPLAIN PLAN FOR + 같은 질의를 쓴 /root/flink/plan/twophase.sql 을 돌려 /root/flink/plan/twophase.out 에 저장하세요. 실행 계획에 GlobalGroupAggregate ← Exchange ← LocalGroupAggregate 와 그 아래 MiniBatchAssigner 가 보여야 합니다.

3단계의 조언이 권한 설정입니다. 두 단계 집계는 mini-batch 가 켜져야 생깁니다 — 섞기 전에 서브태스크마다 미리 합치고(Local), 섞은 뒤 합칩니다(Global). 세 설정 중 하나라도 빠지면 계획이 바뀌지 않습니다.

distinct 분할과 보고서

table.optimizer.distinct-agg.split.enabled = true, table.optimizer.distinct-agg.split.bucket-num = 64 를 SET 하고 EXPLAIN PLAN FOR SELECT page, COUNT(DISTINCT user_id) AS users FROM clicks GROUP BY page; 를 쓴 /root/flink/plan/distinct.sql 의 출력을 /root/flink/plan/distinct.out 에 저장하세요. 그리고 /root/flink/plan/report.jsonjson_nodes·forward_edges·hash_edges(json.out), chained_vertices·unchained_vertices(두 잡 응답), distinct_buckets(distinct.out 의 버킷 수)를 적으세요.

분할이 먹으면 GroupAggregate 가 PARTIAL 과 FINAL 두 층이 되고, 그 사이와 아래에 Exchange 가 하나씩 생기며, 아래 Calc 에 MOD(HASH_CODE(user_id), 버킷 수) 가 보입니다. 보고서 숫자는 전부 여러분이 저장한 파일에서 셀 수 있습니다 — 체이닝 정점 수가 HASH 수 + 1 인지도 확인해 보세요.