Apache Flink — 스트림을 엔진으로 돌린다 · 실행 계획 읽기 · 实验
계획을 읽고 정점 수를 맞힌다
목표
같은 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.sql 의 EXPLAIN PLAN FOR 출력을 /root/flink/plan/explain.out 에 저장하세요.
2. explain.out 의 실행 계획을 읽어 /root/flink/plan/shuffle.json 에 exchange·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 을 쓰세요.
참고
- 원본
/opt/lab/fixtures/data/plan_clicks.csv의 열:click_id BIGINT, user_id STRING, page STRING, amount INT, click_time TIMESTAMP(3)(머리글 없는 CSV). - 질의는 모든 단계가 같습니다:
SELECT page, COUNT(*) AS n_views, SUM(amount) AS revenue FROM clicks WHERE amount > 0 GROUP BY page.views는 예약어라 별칭으로 쓸 수 없습니다. - 표 정의를 매번 되풀이하지 않으려면
sql-client.sh -i ddl.sql -f 파일.sql > 파일.out 2>&1(-i 는 먼저 도는 초기화 파일). sql-client 는 한 줄에 한 문장만 받습니다. - EXPLAIN 출력은 표 한 칸에 여러 줄로 찍힙니다. 실행 계획은 위가 싱크 쪽, 아래가 원천 쪽입니다.
- 흔한 실수:
SET 'execution.runtime-mode' = 'batch'로 돌리면 집계 연산자 이름이 달라집니다 — 이 실습은 기본(스트리밍)으로 합니다. INSERT 가 끝나기 전에/jobs/<jid>를 받으면 RUNNING 입니다(table.dml-sync). - 공식 문서: [EXPLAIN](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/sql/reference/utility/explain/) · [Performance Tuning](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/dev/table/tuning/) · [Table 설정](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/dev/table/config/) · [Flink Architecture](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/concepts/flink-architecture/) · [REST API](https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/ops/rest_api/)
8个步骤
- EXPLAIN 의 세 구역을 받는다
- 섞는 자리를 찾는다
- 비용 추정과 조언을 붙인다
- JSON 실행 계획에서 전달 방식을 본다
- 체이닝을 켠 잡의 정점을 센다
- 체이닝을 끄고 다시 센다
- mini-batch 로 두 단계 집계를 만든다
- distinct 분할과 보고서