SQL 힌트
SQL 힌트 (SQL Hints)
SQL 힌트는 SQL 문과 함께 사용해 실행 계획을 변경할 수 있어요. 이 장은 다양한 접근 방식을 강제하기 위해 힌트를 사용하는 방법을 설명해요.
출처: 문서
본문
배치 | 스트리밍
일반적으로 힌트는 다음에 사용할 수 있어요.
- 플래너 강제 (Enforce planner): 완벽한 플래너는 없으므로, 사용자가 실행을 더 잘 제어할 수 있도록 힌트를 구현하는 것이 합리적이에요.
- 메타데이터(또는 통계) 추가: "스캔용 테이블 인덱스"나 "일부 shuffle 키의 스큐 정보" 같은 일부 통계는 쿼리에 대해 다소 동적이므로, 플래너의 계획 메타데이터가 자주 정확하지 않기 때문에 힌트로 구성하는 것이 매우 편리해요.
- 연산자 리소스 제약: 많은 경우 실행 연산자에 기본 리소스 구성을 제공해요. 예를 들어 최소 parallelism, 관리 메모리(리소스 소모형 UDF), 특별한 리소스 요구사항(GPU 또는 SSD 디스크) 등이에요. 힌트로 Job이 아닌 쿼리별로 리소스를 프로파일링하는 것은 매우 유연해요.
동적 테이블 옵션 (Dynamic Table Options)
동적 테이블 옵션은 SQL DDL이나 connect API로 정의된 정적 테이블 옵션과 달리, 테이블 옵션을 동적으로 지정하거나 덮어쓸 수 있게 해줘요. 이 옵션들은 각 쿼리 내에서 테이블별로 유연하게 지정될 수 있어요.
따라서 대화형 터미널의 ad-hoc 쿼리에 매우 적합해요. 예를 들어 SQL-CLI에서 CSV 소스에 대해 파싱 오류를 무시하도록 지정하려면 동적 옵션 /*+ OPTIONS('csv.ignore-parse-errors'='true') */만 추가하면 돼요.
구문 (Syntax)
SQL 호환성을 깨지 않기 위해 Oracle 스타일 SQL 힌트 구문을 사용해요.
table_path /*+ OPTIONS(key=val [, key=val]*) */
key:
stringLiteral
val:
stringLiteral
예시 (Examples)
CREATE TABLE kafka_table1 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE kafka_table2 (id BIGINT, name STRING, age INT) WITH (...);
-- override table options in query source
select id, name from kafka_table1 /*+ OPTIONS('scan.startup.mode'='earliest-offset') */;
-- override table options in join
select * from
kafka_table1 /*+ OPTIONS('scan.startup.mode'='earliest-offset') */ t1
join
kafka_table2 /*+ OPTIONS('scan.startup.mode'='earliest-offset') */ t2
on t1.id = t2.id;
-- override table options for INSERT target table
insert into kafka_table1 /*+ OPTIONS('sink.partitioner'='round-robin') */ select * from kafka_table2;
쿼리 힌트 (Query Hints)
쿼리 힌트는 지정된 쿼리 범위 내에서 실행 계획에 영향을 주도록 옵티마이저에 제안하는 데 사용할 수 있어요. 그 효과 범위는 Query Hints가 지정된 현재 Query block입니다. 지금은 Flink Query Hints가 Join Hints만 지원해요.
구문 (Syntax)
Flink의 Query Hints 구문은 Apache Calcite의 Query Hints 구문을 따르는 형태예요.
# Query Hints:
SELECT /*+ hint [, hint ] */ ...
hint:
hintName
| hintName '(' optionKey '=' optionVal [, optionKey '=' optionVal ]* ')'
| hintName '(' hintOption [, hintOption ]* ')'
optionKey:
simpleIdentifier
| stringLiteral
optionVal:
stringLiteral
hintOption:
simpleIdentifier
| numericLiteral
| stringLiteral
Query Hints의 충돌 사례
Key-value 힌트 충돌 해결 (Resolution of Key-value Hint Conflicts)
key-value 힌트는 다음 구문으로 제공돼요.
hintName '(' optionKey '=' optionVal [, optionKey '=' optionVal ]* ')'
Flink가 key-value 힌트에서 충돌을 만나면 마지막 쓰기 승리(last-write-wins) 전략을 채택해요. 즉 같은 key에 여러 힌트 값이 제공되면 Flink는 쿼리에서 마지막에 지정된 힌트의 값을 사용해요. 예를 들어 LOOKUP 힌트에서 'max-attempts' 값이 충돌하는 다음 SQL 쿼리를 고려해 봅시다.
SELECT /*+ LOOKUP('table'='D', 'max-attempts'='3', 'max-attempts'='4') */ * FROM t1 T JOIN t2 AS OF T.proctime AS D ON T.id = D.id;
이 경우 Flink는 'max-attempts'에 대해 마지막에 지정된 값을 선택해 충돌을 해결해요. 따라서 'max-attempts'의 유효한 힌트는 '4'가 돼요.
List 힌트 충돌 해결 (Resolution of List Hint Conflicts)
List 힌트는 다음 구문으로 제공돼요.
hintName '(' hintOption [, hintOption ]* ')'
List 힌트에서 Flink는 첫 번째 수락(first-accept) 전략을 채택해 충돌을 해결해요. 즉 목록에서 첫 번째로 지정된 힌트가 우선하고 유효하게 돼요.
예를 들어 BROADCAST 힌트가 충돌하는 다음 SQL 쿼리를 고려해 봅시다.
SELECT /*+ BROADCAST(t2, t1), BROADCAST(t1, t2) */ * FROM t1 JOIN t2 ON t1.id = t2.id;
이 시나리오에서 Flink는 먼저 나열된 BROADCAST 힌트를 선택해요. 따라서 유효한 broadcast 힌트는 BROADCAST(t2, t1)이에요.
조인 힌트 (Join Hints)
Join Hints는 사용자가 더 고성능의 실행 계획을 얻기 위해 옵티마이저에 조인 전략을 제안하도록 허용해요. 지금 Flink Join Hints는 BROADCAST, SHUFFLE_HASH, SHUFFLE_MERGE, NEST_LOOP를 지원해요.
참고:
- Join Hints에 지정된 테이블이 반드시 존재해야 해요. 그렇지 않으면 table not exists 오류가 발생해요.
- Flink Join Hints는 쿼리 블록에서 하나의 힌트 블록만 지원해요.
/*+ BROADCAST(t1) */ /*+ SHUFFLE_HASH(t1) */처럼 여러 힌트 블록이 지정되면 이 쿼리 문을 파싱할 때 예외가 발생해요. - 하나의 힌트 블록에서
/*+ BROADCAST(t1, t2, ..., tn) */처럼 단일 Join Hint에 여러 테이블을 지정하거나,/*+ BROADCAST(t1), BROADCAST(t2), ..., BROADCAST(tn) */처럼 여러 Join Hint를 지정하는 것이 모두 지원돼요. - 단일 Join Hints의 여러 테이블 또는 힌트 블록의 여러 Join Hints에 대해 Flink Join Hints는 충돌할 수 있어요. 충돌이 발생하면 Flink는 가장 잘 맞는 테이블 또는 조인 전략을 선택해요. (Join Hints의 충돌 사례 참고)
BROADCAST
BROADCAST는 Flink가 BroadCast 조인을 사용하도록 제안해요. 힌트가 있는 조인 측은 table.optimizer.join.broadcast-threshold에 관계없이 broadcast되므로, 힌트 측 테이블의 데이터 볼륨이 매우 작을 때 잘 수행돼요.
참고: BROADCAST는 동등 조인 조건(equivalence join condition)이 있는 조인만 지원하며, Full Outer Join은 지원하지 않아요.
예시:
CREATE TABLE t1 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t2 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t3 (id BIGINT, name STRING, age INT) WITH (...);
-- Flink will use broadcast join and t1 will be the broadcast table.
SELECT /*+ BROADCAST(t1) */ * FROM t1 JOIN t2 ON t1.id = t2.id;
-- Flink will use broadcast join for both joins and t1, t3 will be the broadcast table.
SELECT /*+ BROADCAST(t1, t3) */ * FROM t1 JOIN t2 ON t1.id = t2.id JOIN t3 ON t1.id = t3.id;
-- BROADCAST don't support non-equivalent join conditions.
-- Join Hint will not work, and only nested loop join can be applied.
SELECT /*+ BROADCAST(t1) */ * FROM t1 join t2 ON t1.id > t2.id;
-- BROADCAST don't support full outer join.
-- Join Hint will not work in this case, and the planner will choose the appropriate join strategy based on cost.
SELECT /*+ BROADCAST(t1) */ * FROM t1 FULL OUTER JOIN t2 ON t1.id = t2.id;
SHUFFLE_HASH
SHUFFLE_HASH는 Flink가 Shuffle Hash 조인을 사용하도록 제안해요. 힌트가 있는 조인 측이 조인 빌드 측이 되며, 힌트 측 테이블의 데이터 볼륨이 너무 크지 않을 때 잘 수행돼요.
참고: SHUFFLE_HASH는 동등 조인 조건이 있는 조인만 지원해요.
예시:
CREATE TABLE t1 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t2 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t3 (id BIGINT, name STRING, age INT) WITH (...);
-- Flink will use hash join and t1 will be the build side.
SELECT /*+ SHUFFLE_HASH(t1) */ * FROM t1 JOIN t2 ON t1.id = t2.id;
-- Flink will use hash join for both joins and t1, t3 will be the join build side.
SELECT /*+ SHUFFLE_HASH(t1, t3) */ * FROM t1 JOIN t2 ON t1.id = t2.id JOIN t3 ON t1.id = t3.id;
-- SHUFFLE_HASH don't support non-equivalent join conditions.
-- For this case, Join Hint will not work, and only nested loop join can be applied.
SELECT /*+ SHUFFLE_HASH(t1) */ * FROM t1 join t2 ON t1.id > t2.id;
SHUFFLE_MERGE
SHUFFLE_MERGE는 Flink가 Sort Merge 조인을 사용하도록 제안해요. 이 유형의 Join Hint는 두 개의 큰 테이블 간의 조인 시나리오나, 조인 양쪽의 데이터가 이미 정렬된 시나리오에서 사용하는 것을 권장해요.
참고: SHUFFLE_MERGE는 동등 조인 조건이 있는 조인만 지원해요.
예시:
CREATE TABLE t1 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t2 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t3 (id BIGINT, name STRING, age INT) WITH (...);
-- Sort merge join strategy is adopted.
SELECT /*+ SHUFFLE_MERGE(t1) */ * FROM t1 JOIN t2 ON t1.id = t2.id;
-- Sort merge join strategy is both adopted in these two joins.
SELECT /*+ SHUFFLE_MERGE(t1, t3) */ * FROM t1 JOIN t2 ON t1.id = t2.id JOIN t3 ON t1.id = t3.id;
-- SHUFFLE_MERGE don't support non-equivalent join conditions.
-- Join Hint will not work, and only nested loop join can be applied.
SELECT /*+ SHUFFLE_MERGE(t1) */ * FROM t1 join t2 ON t1.id > t2.id;
NEST_LOOP
NEST_LOOP는 Flink가 Nested Loop 조인을 사용하도록 제안해요. 이 유형의 조인 힌트는 특별한 시나리오 요구사항 없이는 권장되지 않아요.
참고: NEST_LOOP는 동등 조인 조건과 비동등 조인 조건 모두를 지원해요.
예시:
CREATE TABLE t1 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t2 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t3 (id BIGINT, name STRING, age INT) WITH (...);
-- Flink will use nested loop join and t1 will be the build side.
SELECT /*+ NEST_LOOP(t1) */ * FROM t1 JOIN t2 ON t1.id = t2.id;
-- Flink will use nested loop join for both joins and t1, t3 will be the join build side.
SELECT /*+ NEST_LOOP(t1, t3) */ * FROM t1 JOIN t2 ON t1.id = t2.id JOIN t3 ON t1.id = t3.id;
MULTI_JOIN
MULTI_JOIN은 Flink가 여러 일반 조인을 동시에 처리하기 위해 MultiJoin 연산자를 사용하도록 제안해요. 하나 이상의 공통 조인 키를 공유하고 큰 중간 상태 또는 레코드 증폭을 경험하는 여러 조인이 있을 때 이 유형의 조인 힌트를 권장해요. MultiJoin 연산자는 다양한 입력 스트림에서 조인을 동시에 처리해 중간 상태를 제거하며, 이는 상태 크기를 크게 줄이고 어떤 경우에는 성능을 향상시킬 수 있어요.
MultiJoin 연산자에 대한 더 자세한 내용(언제 사용해야 하는지, 구성 옵션 포함)은 Multiple Regular Joins를 참고하세요.
참고:
- MULTI_JOIN 힌트는 테이블 이름이나 테이블 별칭을 지정할 수 있어요. 테이블에 별칭이 있으면 힌트는 별칭 이름을 사용해야 해요.
- MultiJoin 연산자를 적용하려면 조인 조건 간에 최소 하나의 키가 공유되어야 해요.
- 지정되면 MULTI_JOIN 힌트는 현재 쿼리 블록 내의 힌트에 나열된 테이블에 적용돼요.
예시:
CREATE TABLE t1 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t2 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t3 (id BIGINT, name STRING, age INT) WITH (...);
-- Flink will use the MultiJoin operator for the three-way join.
SELECT /*+ MULTI_JOIN(t1, t2, t3) */ * FROM t1
JOIN t2 ON t1.id = t2.id
JOIN t3 ON t1.id = t3.id;
-- Using table names instead of aliases.
SELECT /*+ MULTI_JOIN(Users, Orders, Payments) */ * FROM Users
INNER JOIN Orders ON Users.user_id = Orders.user_id
INNER JOIN Payments ON Users.user_id = Payments.user_id;
-- Partial match: only t1 and t2 will use MultiJoin, t3 will use regular join.
SELECT /*+ MULTI_JOIN(t1, t2) */ * FROM t1
JOIN t2 ON t1.id = t2.id
JOIN t3 ON t1.id = t3.id;
-- Combining MULTI_JOIN with STATE_TTL hint.
SELECT /*+ MULTI_JOIN(t1, t2, t3), STATE_TTL('t1'='1d', 't2'='2d', 't3'='12h') */ * FROM t1
JOIN t2 ON t1.id = t2.id
JOIN t3 ON t1.id = t3.id;
LOOKUP
LOOKUP 힌트는 사용자가 Flink 옵티마이저에게 다음을 제안하도록 허용해요.
- 동기(sync) 또는 비동기(async) 조회 함수 사용
- async 매개변수 구성
- 조회를 위한 지연 재시도(delayed retry) 전략 활성화
LOOKUP 힌트 옵션:
| option type | option name | required | value type | default value | description |
|---|---|---|---|---|---|
| table | table | Y | string | N/A | 조회 소스의 테이블 이름 |
| async | async | N | boolean | N/A | 값은 'true' 또는 'false'일 수 있으며, 플래너가 해당 조회 함수를 선택하도록 제안해요. 백엔드 조회 소스가 제안된 조회 모드를 지원하지 않으면 효과가 없어요. |
| output-mode | N | string | ordered | 값은 'ordered' 또는 'allow_unordered'일 수 있어요. 'allow_unordered'는 사용자가 순서 없는 결과를 허용하면, 결과의 정확성에 영향을 주지 않을 때 AsyncDataStream.OutputMode.UNORDERED를 시도한다는 뜻이에요. 그렇지 않으면 여전히 ORDERED가 사용돼요. 이는 ExecutionConfigOptions#TABLE_EXEC_ASYNC_LOOKUP_OUTPUT_MODE와 일치해요. |
|
| capacity | N | integer | 100 | 조회 조인 연산자의 백엔드 asyncWaitOperator를 위한 버퍼 용량 | |
| timeout | N | duration | 300s | 비동기 연산의 첫 호출부터 최종 완료까지의 타임아웃. 여러 번의 재시도를 포함할 수 있으며, 페일오버 시 재설정돼요 | |
| retry | retry-predicate | N | string | N/A | 'lookup_miss'일 수 있으며, 조회 결과가 비어 있으면 재시도를 활성화해요. |
| retry-strategy | N | string | N/A | 'fixed_delay'일 수 있어요. | |
| fixed-delay | N | duration | N/A | 'fixed_delay' 전략을 위한 지연 시간 | |
| max-attempts | N | integer | N/A | 'fixed_delay' 전략의 최대 재시도 횟수 | |
| shuffle | shuffle | N | boolean | false | 사용자 정의 조회 shuffle을 활성화할지 여부. 조회 소스가 입력 데이터 분배를 결정하고 그에 따라 조회 전략을 최적화할 수 있게 해요 |
참고:
- 'table' 옵션은 필수이며, 테이블 이름만 지원돼요 (FROM 절과 일치하게 유지). 테이블에 별칭 이름이 있으면 별칭 이름만 사용할 수 있어요.
- async 옵션은 모두 선택 사항이며, 구성하지 않으면 기본값을 사용해요.
- retry 옵션에는 기본값이 없어서, 재시도를 활성화해야 할 때 모든 retry 옵션을 유효한 값으로 설정해야 해요.
1. Sync 및 Async 조회 함수 사용
커넥터가 async와 sync 조회 능력을 모두 가지고 있다면, 사용자는 'async'='false' 옵션 값으로 플래너가 sync 조회를 사용하도록 제안하거나 'async'='true'로 async 조회를 사용하도록 제안할 수 있어요.
예시:
-- suggest the optimizer to use sync lookup
LOOKUP('table'='Customers', 'async'='false')
-- suggest the optimizer to use async lookup
LOOKUP('table'='Customers', 'async'='true')
참고: 'async' 옵션이 지정되지 않으면 옵티마이저는 async 조회를 선호해요. 다음의 경우 항상 sync 조회를 사용해요.
- 커넥터가 sync 조회만 구현했을 때
- 사용자가 'table.optimizer.non-deterministic-update.strategy'의 'TRY_RESOLVE' 모드를 활성화했고 옵티마이저가 비결정적 업데이트로 인한 정확성 문제를 확인했을 때
2. Async 매개변수 구성
사용자는 async 조회 모드에서 async 옵션을 통해 async 매개변수를 구성할 수 있어요.
예시:
-- configure the async parameters: 'output-mode', 'capacity', 'timeout', can set single one or multi params
LOOKUP('table'='Customers', 'async'='true', 'output-mode'='allow_unordered', 'capacity'='100', 'timeout'='180s')
예를 들어 job 수준 구성이 다음과 같다면:
table.exec.async-lookup.output-mode: ORDERED
table.exec.async-lookup.buffer-capacity: 100
table.exec.async-lookup.timeout: 180s
다음 힌트들은:
1. LOOKUP('table'='Customers', 'async'='true', 'output-mode'='allow_unordered')
2. LOOKUP('table'='Customers', 'async'='true', 'timeout'='300s')
다음과 동일해요:
1. LOOKUP('table'='Customers', 'async'='true', 'output-mode'='allow_unordered', 'capacity'='100', 'timeout'='180s')
2. LOOKUP('table'='Customers', 'async'='true', 'output-mode'='ordered', 'capacity'='100', 'timeout'='300s')
3. 조회를 위한 지연 재시도 전략 활성화
조회 조인의 지연 재시도는 외부 시스템의 지연 업데이트가 스트림 데이터와의 예상치 못한 보강(enrichment)을 일으키는 문제를 해결하기 위한 것이에요. 힌트 옵션 'retry-predicate'='lookup_miss'은 sync와 async 조회 모두에서 재시도를 활성화할 수 있으며, 현재는 fixed delay 재시도 전략만 지원돼요.
fixed delay 재시도 전략의 옵션:
'retry-strategy'='fixed_delay'
-- fixed delay duration
'fixed-delay'='10s'
-- max number of retry(counting from the retry operation, if set to '1', then a single lookup process
-- for a specific lookup key will actually execute up to 2 lookup requests)
'max-attempts'='3'
예시:
- async 조회에서 재시도 활성화
LOOKUP('table'='Customers', 'async'='true', 'retry-predicate'='lookup_miss', 'retry-strategy'='fixed_delay', 'fixed-delay'='10s','max-attempts'='3')
- sync 조회에서 재시도 활성화
LOOKUP('table'='Customers', 'async'='false', 'retry-predicate'='lookup_miss', 'retry-strategy'='fixed_delay', 'fixed-delay'='10s','max-attempts'='3')
조회 소스가 하나의 능력만 가진다면 'async' 모드 옵션은 생략할 수 있어요.
LOOKUP('table'='Customers', 'retry-predicate'='lookup_miss', 'retry-strategy'='fixed_delay', 'fixed-delay'='10s','max-attempts'='3')
4. 사용자 정의 데이터 분배 활성화
기본적으로 Lookup Join의 입력 스트림 데이터 배포는 임의적이므로, 소스가 캐시를 효과적으로 활용해 조회를 가속화하지 못할 수 있어요. 다음과 같이 custom shuffle을 활성화하면 소스가 입력 데이터의 배포를 스스로 결정할 수 있고, 이 사전 지식을 사용해 캐시와 조회 전략을 최적화할 수 있어요.
LOOKUP('table'='Customers', 'shuffle'='true')
이 기능을 완전히 활용하려면 대상 조회 소스가 custom shuffle을 지원해야 해요. 커넥터 개발자에게는 LookupTableSource 서브클래스가 SupportsLookupCustomShuffle을 구현하면 이를 달성할 수 있어요. 소스가 아직 그러한 지원을 제공하지 않더라도 사용자는 먼저 이 기능을 활성화할 수 있고, Flink가 최선을 다해 해시 파티셔닝을 적용하며 이것도 성능 향상을 가져와야 해요.
추가 참고 사항
재시도에 대한 캐싱 활성화의 효과
FLIP-221은 조회 소스에 캐싱 지원을 추가하며, 여기에는 PARTIAL과 FULL 캐싱 모드가 있어요 (모드 NONE은 캐싱 비활성화를 의미). FULL 캐싱이 활성화되면 재시도가 전혀 없어요 (조회 소스의 전체 캐시된 미러를 통해 조회를 재시도하는 것은 무의미하기 때문). PARTIAL 캐싱이 활성화되면 들어오는 레코드에 대해 먼저 로컬 캐시에서 조회하고, 캐시 미스 시 백엔드 커넥터를 통해 외부 조회를 해요 (캐시 히트 시 레코드를 즉시 반환). 그리고 조회 결과가 비면 (캐싱 비활성화와 동일하게) 재시도를 트리거하며, 최종 조회 결과는 재시도 완료 시 결정돼요 (PARTIAL 캐싱 모드에서도 로컬 캐시를 업데이트해요).
조회 키와 'retry-predicate'='lookup_miss' 재시도 조건에 대한 참고
커넥터마다 인덱스 조회 능력이 다를 수 있어요. 예를 들어 내장 HBase 커넥터는 rowkey로만 (보조 인덱스 없이) 조회할 수 있는 반면, 내장 JDBC 커넥터는 임의의 컬럼에 대해 더 강력한 인덱스 조회 능력을 제공할 수 있어요. 이는 서로 다른 물리적 저장소에 의해 결정돼요.
여기서 말하는 조회 키는 인덱스 조회를 위한 필드 또는 필드의 조합이에요. 조인 조인의 예시처럼 c.id가 조인 조건 "ON o.customer_id = c.id"의 조회 키예요.
SELECT o.order_id, o.total, c.country, c.zip
FROM Orders AS o
JOIN Customers FOR SYSTEM_TIME AS OF o.proc_time AS c
ON o.customer_id = c.id
조인 조건을 "ON o.customer_id = c.id and c.country = 'US'"로 바꾼다면:
SELECT o.order_id, o.total, c.country, c.zip
FROM Orders AS o
JOIN Customers FOR SYSTEM_TIME AS OF o.proc_time AS c
ON o.customer_id = c.id and c.country = 'US'
Customers 테이블이 MySql에 저장되어 있다면 c.id와 c.country 모두 조회 키로 사용돼요.
CREATE TEMPORARY TABLE Customers (
id INT,
name STRING,
country STRING,
zip STRING
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://mysqlhost:3306/customerdb',
'table-name' = 'customers'
)
Customers 테이블이 HBase에 저장되어 있다면 c.id만 조회 키가 될 수 있고, 나머지 조인 조건 c.country = 'US'는 조회 결과가 반환된 후 평가돼요.
CREATE TEMPORARY TABLE Customers (
id INT,
name STRING,
country STRING,
zip STRING,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'hbase-2.2',
...
)
따라서 위 쿼리는 'lookup_miss' 재시도 술어와 fixed-delay 재시도 전략을 활성화할 때 서로 다른 저장소에서 서로 다른 재시도 효과를 가져요.
예를 들어 Customers 테이블에 다음과 같은 행이 있다면:
id=100, country='CN'
order 스트림에서 'id=100'인 레코드를 처리할 때, 'jdbc' 커넥터에서는 c.id와 c.country 모두 조회 키로 사용되므로 (country='CN'은 조건 c.country = 'US'를 만족하지 않으므로) 해당 조회 결과는 null이 되고, 이는 재시도를 트리거해요. 'hbase' 커넥터에서는 c.id만 조회 키로 사용되므로 해당 조회 결과는 비어있지 않게 되고 (레코드 id=100, country='CN'을 반환), 이는 재시도를 트리거하지 않아요 (나머지 조인 조건 c.country = 'US'는 반환된 레코드에 대해 참이 아닌 것으로 평가됨).
현재 SQL 의미론 고려에 기반해 'lookup_miss' 재시도 술어만 제공돼요. 그리고 (테이블에 없는 것이 아니라 이미 존재하는 과거 버전 레코드인) 차원 테이블의 지연 업데이트를 기다리는 것이 필요할 때, 사용자는 두 가지 해결책을 시도할 수 있어요.
- DataStream Async I/O의 새 재시도 지원으로 사용자 정의 재시도 술어를 구현 (반환된 레코드에 대해 더 복잡한 판단 허용)
- 위 예시의 경우 타임스탬프로 생성된 어떤 종류의 데이터 버전에 대한 비교를 포함하는 다른 조인 조건을 추가해 지연 재시도를 활성화. Customers 테이블이 매시간 업데이트된다고 가정하면, 시간별 정밀도로 예약된 새 시간 종속 버전 필드
update_version을 추가할 수 있어요. 예를 들어 레코드의 업데이트 시간 '2022-08-15 12:01:02'는 update_version을 '2022-08-15 12:00'으로 저장해요.
CREATE TEMPORARY TABLE Customers (
id INT,
name STRING,
country STRING,
zip STRING,
-- the newly added time-dependent version field
update_version STRING
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://mysqlhost:3306/customerdb',
'table-name' = 'customers'
)
조인 조건에 Order 스트림의 시간 필드와 Customers.update_version 양쪽에 대한 동등 조건을 추가해요.
ON o.customer_id = c.id AND DATE_FORMAT(o.order_timestamp, 'yyyy-MM-dd HH:mm') = c.update_version
그러면 Order의 레코드가 Customers 테이블에서 '12:00' 버전의 새 레코드를 조회하지 못할 때 지연 재시도를 활성화할 수 있어요.
문제 해결 (Trouble Shooting)
지연 재시도 조회를 켜면 조회 노드에서 백프레셔(backpressure) 문제를 만날 가능성이 더 높아요. 이는 web ui의 'Task Manager' 페이지에서 'Thread Dump'를 통해 빠르게 확인할 수 있어요. async 및 sync 조회에서 각각, thread sleep의 호출 스택이 나타나요.
- async 조회:
RetryableAsyncLookupFunctionDelegator - sync 조회:
RetryableLookupFunctionDelegator
참고:
- 재시도가 있는 async 조회는 모든 입력 데이터에 대해 고정된 지연 처리가 불가능해요 (다른 더 가벼운 방식으로 해결해야 해요. 예: pending source consumption 또는 재시도가 있는 sync 조회 사용)
- sync 조회에서 재시도 실행을 위한 지연 대기는 완전히 동기적이에요. 즉 현재 레코드가 완료될 때까지 다음 레코드의 처리가 시작되지 않아요.
- async 조회에서 'output-mode'가 'ORDERED' 모드이면 지연 재시도로 인한 백프레셔 확률이 'UNORDERED' 모드보다 더 높을 수 있어요. 이 경우 async 'capacity'를 늘리는 것이 백프레셔 감소에 효과적이지 않을 수 있으며, 지연 기간을 줄이는 것을 고려해야 할 수 있어요.
Join Hints의 충돌 사례 (Conflict Cases In Join Hints)
Join Hints가 충돌하면 Flink는 가장 잘 맞는 것을 선택해요.
- 먼저, Join Hints는 충돌 해결을 위해 Flink query hint의 논리를 따르는 형태예요 (Query Hints의 충돌 사례 참고)
- 같은 Join Hint 전략에서의 충돌: Flink는 조인에 대해 첫 번째로 일치하는 테이블을 선택해요.
- 다른 Join Hints 전략 간의 충돌: Flink는 조인에 대해 첫 번째로 일치하는 힌트를 선택해요.
예시:
CREATE TABLE t1 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t2 (id BIGINT, name STRING, age INT) WITH (...);
CREATE TABLE t3 (id BIGINT, name STRING, age INT) WITH (...);
-- Conflict in One Same Join Hints Strategy Case
-- The first hint will be chosen, t2 will be the broadcast table.
SELECT /*+ BROADCAST(t2), BROADCAST(t1) */ * FROM t1 JOIN t2 ON t1.id = t2.id;
-- BROADCAST(t2, t1) will be chosen, and t2 will be the broadcast table.
SELECT /*+ BROADCAST(t2, t1), BROADCAST(t1, t2) */ * FROM t1 JOIN t2 ON t1.id = t2.id;
-- This case equals to BROADCAST(t1, t2) + BROADCAST(t3),
-- when join between t1 and t2, t1 will be the broadcast table,
-- when join between the result after t1 join t2 and t3, t3 will be the broadcast table.
SELECT /*+ BROADCAST(t1, t2, t3) */ * FROM t1 JOIN t2 ON t1.id = t2.id JOIN t3 ON t1.id = t3.id;
-- Conflict in Different Join Hints Strategies Case
-- The first Join Hint (BROADCAST(t1)) will be chosen, and t1 will be the broadcast table.
SELECT /*+ BROADCAST(t1) SHUFFLE_HASH(t1) */ * FROM t1 JOIN t2 ON t1.id = t2.id;
-- Although BROADCAST is first one hint, but it doesn't support full outer join,
-- so the following SHUFFLE_HASH(t1) will be chosen, and t1 will be the join build side.
SELECT /*+ BROADCAST(t1) SHUFFLE_HASH(t1) */ * FROM t1 FULL OUTER JOIN t2 ON t1.id = t2.id;
-- Although there are two Join Hints were defined, but all of them are neither support non-equivalent join,
-- so only nested loop join can be applied.
SELECT /*+ BROADCAST(t1) SHUFFLE_HASH(t1) */ * FROM t1 FULL OUTER JOIN t2 ON t1.id > t2.id;
State TTL 힌트 (State TTL Hints)
상태 저장 연산인 Regular Join과 Group Aggregation에 대해, 사용자는 STATE_TTL 힌트로 연산자 수준의 Idle State Retention Time을 지정할 수 있어요. 이는 앞서 언급된 연산자가 파이프라인 수준 구성 table.exec.state.ttl과 다른 TTL을 갖게 해요.
Regular Join 예시
CREATE TABLE orders (
o_orderkey INT,
o_custkey INT,
o_status BOOLEAN,
o_totalprice DOUBLE
) WITH (...);
CREATE TABLE lineitem (
l_linenumber int,
l_orderkey int,
l_partkey int,
l_extendedprice double
) WITH (...);
CREATE TABLE customers (
c_custkey int,
c_address string
) WITH (...);
-- table name as hint key
SELECT /*+ STATE_TTL('orders'='3d', 'lineitem'='1d') */ * FROM
orders LEFT JOIN lineitem
ON orders.o_orderkey = lineitem.l_orderkey;
-- table alias as hint key
SELECT /*+ STATE_TTL('o'='3d', 'l'='1d') */ * FROM
orders o LEFT JOIN lineitem l
ON o.o_orderkey = l.l_orderkey;
-- temporary view name as hint key
CREATE TEMPORARY VIEW left_input AS SELECT ... FROM orders WHERE ...;
CREATE TEMPORARY VIEW right_input AS SELECT ... FROM lineitem WHERE ...;
SELECT /*+ STATE_TTL('left_input'= '360000s', 'right_input' = '15h') */ *
FROM left_input JOIN right_input
ON left_input.join_key = right_input.join_key;
-- cascade joins
SELECT /*+ STATE_TTL('o' = '3d', 'l' = '1d', 'c' = '10d') */ *
FROM orders o LEFT OUTER JOIN lineitem l
ON o.o_orderkey = l.l_orderkey
LEFT OUTER JOIN customers c
ON o.o_custkey = c.c_custkey;
Group Aggregation 예시
-- table name as hint key
SELECT /*+ STATE_TTL('orders' = '1d') */ o_orderkey, SUM(o_totalprice) AS revenue
FROM orders
GROUP BY o_orderkey;
-- table alias as hint key
SELECT /*+ STATE_TTL('o' = '1d') */ o_orderkey, SUM(o_totalprice) AS revenue
FROM orders AS o
GROUP BY o_orderkey;
-- query block alias as hint key
SELECT /*+ STATE_TTL('tmp' = '1d') */ o_orderkey, SUM(o_totalprice) AS revenue
FROM (SELECT o_orderkey, o_totalprice
FROM orders
WHERE o_shippriority = 0) tmp
GROUP BY o_orderkey;
참고:
- 사용자는 테이블/뷰 이름 또는 테이블 별칭을 힌트 키로 선택할 수 있어요. 단, 별칭이 지정되면 STATE_TTL은 별칭에 힌트되어야 해요.
CREATE TEMPORARY VIEW V AS
SELECT /*+ STATE_TTL('A' = '1d', 'B' = '12h')*/ * FROM A JOIN B ON...;
SELECT /*+ STATE_TTL('V' = '1d', 'C' = '3d')*/ * FROM V JOIN C ON ...;
- cascade join의 경우, 지정된 state TTL은 (상향식 순서로) 첫 번째 조인 연산자의 좌·우 state TTL과 두 번째 조인 연산자의 우 state TTL로 해석돼요. 두 번째 조인 연산자의 좌 state TTL은 구성
table.exec.state.ttl에서 가져와요. 두 번째 조인 연산자의 좌 상태에 특정 TTL 값을 설정해야 한다면 쿼리를 쿼리 블록으로 나눠야 해요. - STATE_TTL 힌트는 기본 쿼리 블록에만 적용돼요.
- STATE_TTL 힌트 키가 중복되면 마지막 발생에서 값이 적용돼요. 예를 들어
SELECT /*+ STATE_TTL('A' = '1d', 'A' = '2d')*/ * FROM ...같은 경우 입력 A의 TTL은 2d로 간주돼요. - 중복된 힌트 키가 있는 여러 STATE_TTL 힌트가 나타나면 첫 번째 발생에서 값이 적용돼요. 예를 들어
SELECT /*+ STATE_TTL('A' = '1d', 'B' = '2d'), STATE_TTL('C' = '12h', 'A' = '6h')*/ * FROM ...같은 경우 입력 A의 TTL은 1d로 간주돼요.
쿼리 블록이란 무엇인가? (What are query blocks ?)
쿼리 블록은 SQL의 기본 단위예요. 예를 들어 어떤 인라인 뷰나 하위 쿼리도 외부 쿼리에 대해 별도의 쿼리 블록으로 간주돼요.
예시
SQL 문은 여러 하위 쿼리로 구성될 수 있어요. 하위 쿼리는 SELECT, INSERT 또는 DELETE일 수 있어요. 하위 쿼리는 FROM 절, WHERE 절, 또는 UNION/UNION ALL의 하위 SELECT에 다른 하위 쿼리를 포함할 수 있어요.
이러한 서로 다른 하위 쿼리 또는 뷰 유형에 대해 여러 쿼리 블록으로 구성될 수 있어요. 예를 들어:
아래의 단순한 쿼리는 하위 쿼리가 하나뿐이지만, 두 개의 쿼리 블록이 있어요 — 하나는 외부 SELECT용이고 다른 하나는 하위 쿼리 SELECT용이에요.
아래의 쿼리는 union 쿼리이며, 두 개의 쿼리 블록이 포함돼요 — 하나는 첫 번째 SELECT용이고 다른 하나는 두 번째 SELECT용이에요.
아래의 쿼리는 뷰를 포함하며, 두 개의 쿼리 블록이 있어요 — 하나는 외부 SELECT용이고 다른 하나는 뷰용이에요.