Apache Flink — Running Streams on a Real Engine
Read the plan and predict the vertex count
한국어 원문으로 표시합니다.
목표
같은 GROUP BY 질의의 실행 계획을 EXPLAIN 으로 받아 연산자와 Exchange 의 자리를 읽고, 실제 잡의 정점 수를 체이닝과 전달 방식으로 설명한다. mini-batch 두 단계 집계와 distinct 분할이 계획에 어떻게 나타나는지 확인한다.
왜 중요한가
느린 잡의 원인 — 원천까지 내려가지 못한 필터, 한 키에 몰린 집계, 끊긴 체인 — 은 처리량 그래프에서는 똑같이 보이고 계획에서는 다르게 보인다. 튜닝 설정도 계획이 바뀌어야 먹은 것이다. 이 실습의 채점기는 클러스터에 묻지 않는다. 여러분이 저장한 EXPLAIN 출력 원문과 REST 응답만 읽고, 정점 수는 JSON 실행 계획의 전달 방식과 교차 대조한다.
단계
flink-up뒤 /root/flink/plan/ddl.sql 에 원천clicks와 blackhole 싱크page_stats를 만들고, /root/flink/plan/explain.sql 의EXPLAIN PLAN FOR출력을 /root/flink/plan/explain.out 에 저장하세요.- explain.out 의 실행 계획을 읽어 /root/flink/plan/shuffle.json 에
exchange·above·below·filter_pushed_into_scan을 적으세요. EXPLAIN ESTIMATED_COST, PLAN_ADVICE를 쓴 /root/flink/plan/advice.sql 의 출력을 /root/flink/plan/advice.out 에 저장하세요.EXPLAIN JSON_EXECUTION_PLAN INSERT INTO page_stats …를 쓴 /root/flink/plan/json.sql 의 출력을 /root/flink/plan/json.out 에 저장하세요.- 잡 이름
flk-plan-chained로 같은 INSERT 를 끝까지 돌리는 /root/flink/plan/chained.sql 을 실행하고 그 잡의/jobs/<jid>를 /root/flink/plan/chained-job.json 에 저장하세요. - 체이닝을 끄고 잡 이름
flk-plan-unchained로 돌리는 /root/flink/plan/unchained.sql 을 실행하고/jobs/<jid>를 /root/flink/plan/unchained-job.json 에 저장하세요. - mini-batch 세 설정을 준 뒤 같은 SELECT 를 EXPLAIN 하는 /root/flink/plan/twophase.sql 의 출력을 /root/flink/plan/twophase.out 에 저장하세요.
- 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 · Performance Tuning · Table 설정 · Flink Architecture · REST API
EXPLAIN 의 세 구역을 받는다
flink-up 뒤 /root/flink/plan/ddl.sql 에 plan_clicks.csv 를 읽는 clicks 와 CREATE TABLE page_stats (page STRING, n_views BIGINT, revenue BIGINT) WITH ('connector' = 'blackhole') 를 쓰고, /root/flink/plan/explain.sql 에 EXPLAIN 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.json 에 exchange(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.json 에 json_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 인지도 확인해 보세요.