물리적 옵티마이저
물리적 옵티마이저 (Physical Optimizer)
새로운 Multistage Engine Physical Query Optimizer를 설명해요.
참고: Physical Optimizer는 multi-stage 엔진의 선택적 쿼리 옵티마이저로, Pinot 1.4에서 도입되었고 Pinot 1.5.0부터 안정화되었어요.
출처: 문서
본문
Multistage Engine에 새 쿼리 옵티마이저를 추가했어요. 이 옵티마이저는 Sort Pushdown, Aggregate Split/Pushdown 같은 중요한 최적화를 실행하기 전에 전체 계획에 걸쳐 정확한 Data Distribution을 계산하고 추적해요.
이 옵티마이저의 가장 큰 기능 중 하나는 임의로 복잡한 쿼리에 대해 Query Hints 없이도 가능할 때 Shuffles를 제거하거나 Exchanges를 단순화할 수 있다는 점이에요.
MSE 쿼리에 대해 이 옵티마이저를 활성화하려면 다음 Query Options를 사용할 수 있어요.
SET useMultistageEngine=true;
SET usePhysicalOptimizer=true;
핵심 기능
아래 예시는 COLOCATED_JOIN Quickstart를 기반으로 해요.
자동 Colocated 조인 및 셔플 단순화
3개의 조인으로 구성된 아래 쿼리를 고려해 보세요. 새 쿼리 옵티마이저를 사용하면 데이터가 userUUID로 호환 가능한 수의 파티션으로 분할되어 있으므로(아래 "Setting Up Table Data Distribution" 섹션 참고), 전체 쿼리가 서버 간 데이터 교환 없이 실행될 수 있어요.
SET useMultistageEngine = true;
SET usePhysicalOptimizer = true;
WITH filtered_users AS (
SELECT
userUUID
FROM userAttributes
WHERE userUUID NOT IN (
SELECT
userUUID
FROM userGroups
WHERE groupUUID = 'group-1'
)
AND userUUID IN (
SELECT
userUUID
FROM userGroups
WHERE groupUUID = 'group-2'
)
)
SELECT
userUUID,
SUM(tripAmount)
FROM userFactEvents
WHERE
userUUID IN (
SELECT userUUID FROM filtered_users
)
GROUP BY userUUID
이 쿼리의 쿼리 계획은 아래와 같아요. 전체 쿼리가 아래 Exchange Types에 정의된 1:1 Exchange인 IDENTITY_EXCHANGE를 활용하는 것을 볼 수 있어요.
PhysicalExchange(exchangeStrategy=[SINGLETON_EXCHANGE])
PhysicalAggregate(group=[{1}], agg#0=[$SUM0($0)], aggType=[DIRECT])
PhysicalJoin(condition=[=($1, $2)], joinType=[semi])
PhysicalExchange(exchangeStrategy=[IDENTITY_EXCHANGE])
PhysicalProject(tripAmount=[$7], userUUID=[$10])
PhysicalTableScan(table=[[default, userFactEvents]])
PhysicalJoin(condition=[=($0, $1)], joinType=[semi])
PhysicalProject(userUUID=[$0])
PhysicalFilter(condition=[IS NOT TRUE($3)])
PhysicalJoin(condition=[=($1, $2)], joinType=[left])
PhysicalExchange(exchangeStrategy=[IDENTITY_EXCHANGE])
PhysicalProject(userUUID=[$6], userUUID0=[$6])
PhysicalTableScan(table=[[default, userAttributes]])
PhysicalExchange(exchangeStrategy=[IDENTITY_EXCHANGE])
PhysicalAggregate(group=[{0}], agg#0=[MIN($1)], aggType=[DIRECT])
PhysicalProject(userUUID=[$4], $f1=[true])
PhysicalFilter(condition=[=($3, _UTF-8'group-1')])
PhysicalTableScan(table=[[default, userGroups]])
PhysicalExchange(exchangeStrategy=[IDENTITY_EXCHANGE])
PhysicalProject(userUUID=[$4])
PhysicalFilter(condition=[=($3, _UTF-8'group-2')])
PhysicalTableScan(table=[[default, userGroups]])
다른 서버/파티션 수로 셔플 단순화
새 옵티마이저는 다음 경우에도 셔플을 단순화할 수 있어요.
- 조인의 한쪽이 사용하는 서버가 다를 때
- 조인 입력의 파티션 수가 다를 때
아래 예시에서 orange(왼쪽)와 green(오른쪽) 두 테이블에 걸친 조인을 수행해요. orange 테이블은 4개 파티션이 있고 green 테이블은 2개 파티션이 있어요. orange와 green 테이블에 선택된 서버는 각각 [S0, S1]과 [S0, S2]예요. Physical Optimizer는 기본적으로 왼쪽-가장 입력 연산자와 같은 Workers를 사용하므로 조인은 서버 [S0, S1]에서 수행돼요.
두 테이블을 분할하는 데 사용되는 해시 함수가 같다면 Identity Exchange를 활용하고 조인 양쪽에서 데이터를 재분할하는 것을 건너뛸 수 있어요. 이는 S0이 orange 테이블의 파티션 P_0과 P_2의 레코드로 구성되고, 둘을 합치면 2로 나눈 파티션 P_0을 구성하는 모든 레코드를 포함하기 때문이에요. 즉:
{(P_0 ∪ P_2)}*{mod 4} = (P_0)*{mod 2}
Identity Exchange가 발신자와 수신자의 서버가 같음을 의미하지는 않는다는 점에 주의해요. 이는 발신자에서 수신자로 1:1 매핑이 있음을 의미할 뿐이에요. 아래 예시에서 S2에서 S1으로의 데이터 전송은 네트워크를 통해 이루어져요.

집계 Exchange 자동 생략
GROUP BY userUUID 같은 것을 정확히 평가하려면 userUUID 컬럼을 기준으로 레코드를 분산해야 해요. 기존 쿼리 옵티마이저는 is_partitioned_by_group_by_keys 쿼리 힌트를 사용하지 않는 한 각 집계 아래에 Partitioning Exchange를 추가했어요.
Physical Optimizer는 데이터가 이미 필요한 컬럼으로 분할되어 있는지 감지하고 Exchange 추가를 자동으로 건너뛸 수 있어요. 여기에는 두 가지 이점이 있어요.
- 불필요한 Data Exchanges를 피해요.
- 기본적으로 Exchange 위에 Aggregate가 있으면 Exchange 아래에 Aggregate 복사본이 추가되므로(
is_skip_leaf_stage_group_by쿼리 힌트가 설정되지 않은 경우) Aggregate 분할을 피해요.
이 최적화는 위에 공유한 쿼리 예시에서 실제로 볼 수 있어요. 데이터가 이미 userUUID로 분할되어 있으므로 모든 집계가 DIRECT 모드로 실행되며, 즉 집계를 여러 개로 분할하지 않고 실행돼요.
세그먼트 / 서버 프루닝
Single Stage Engine과 유사하게 테이블의 Routing 구성에서 segmentPrunerTypes를 활성화했다면 Physical Optimizer는 잎 스테이지에 대해 시간, 파티션 또는 다른 프루너 타입을 사용해 세그먼트와 서버를 프루닝해요. 예를 들어 다음 쿼리는 다음 제약 조건을 만족하는 세그먼트만 선택해요.
segmentPartition = Murmur("user-1") % numPartitions
SET useMultistageEngine = true;
SET usePhysicalOptimizer = true;
WITH user_events AS (
SELECT
productCode, tripAmount
FROM
userFactEvents
WHERE
userUUID = 'user-1'
ORDER BY
ts
DESC
LIMIT 100
)
SELECT
productCode,
SUM(tripAmount)
FROM
user_events
GROUP BY productCode
분할이 주어진 파티션에 해당하는 세그먼트가 오직 1개 서버에만 존재하도록 이루어진다면 위 전체 쿼리는 단일 서버 내에서 실행되어 다른 시스템의 shard-local 실행을 시뮬레이션해요.
Pinot Broker에서 상수 쿼리 해결
Apache Calcite는 항상 False로 평가되는 Filter Expression을 감지할 수 있어요. 이러한 경우 쿼리 계획에 Table Scan이 전혀 없을 수 있어요. Physical Optimizer는 그러한 쿼리를 서버를 전혀 개입시키지 않고 Broker 내에서 해결해요.
SET useMultistageEngine = true;
SET usePhysicalOptimizer = true;
SELECT
COUNT(*)
FROM
userFactEvents
WHERE
userUUID = 'user-1' AND userUUID = 'user-2'
Worker 할당
현재 Worker 할당은 다음의 간단한 규칙을 따르는 단계를 거쳐요.
- 잎 스테이지는 Table Scan과 Filters를 기반으로 Table Config에 설정된 Routing 구성을 사용해 workers가 할당돼요.
- 다른 스테이지는 왼쪽-가장 입력 스테이지와 같은 workers를 사용해요.
Sort(fetch=..)같은 일부 Plan Nodes는 데이터를 단일 Worker로 수집해야 할 수 있어요. 그러한 경우 해당 스테이지는 입력 workers 중 하나에서 무작위로 선택된 단일 Worker에서 실행돼요.
제한 사항
기존 MSE 쿼리 옵티마이저의 기능 중 일부는 아직 Physical Optimizer에서 사용할 수 없어요. 대부분은 Pinot 1.5에서 지원을 추가할 계획이에요.
- Spools
- semi-join을 위한 동적 필터