SQL 힌트
SQL 힌트 (SQL Hints)
SQL 힌트는 실행 계획을 변경하기 위해 SQL 문과 함께 사용될 수 있습니다. 이 장에서는 다양한 방식을 강제하기 위해 힌트를 사용하는 방법을 설명합니다.
출처: 문서
본문
이 절은 Batch 와 Streaming 모드에서 모두 지원됩니다.
SQL 힌트는 실행 계획을 변경하기 위해 SQL 문과 함께 사용될 수 있습니다. 이 장에서는 다양한 방식을 강제하기 위해 힌트를 사용하는 방법을 설명합니다.
일반적으로 힌트는 다음에 사용될 수 있습니다:
- 플래너 강제(Enforce planner): 완벽한 플래너는 없으므로 사용자가 실행을 더 잘 제어할 수 있도록 힌트를 구현하는 것이 의미가 있습니다.
- 메타데이터(또는 통계) 추가: "스캔용 테이블 인덱스" 나 "일부 셔플 키의 왜곡 정보" 같은 일부 통계는 쿼리에 대해 다소 동적이므로, 플래너의 계획 메타데이터가 종종 그리 정확하지 않기 때문에 힌트로 구성하는 것이 매우 편리합니다.
- operator 리소스 제약: 많은 경우 실행 operator 에 대해 기본 리소스 구성을 제공합니다. 즉 최소 병렬도 또는 관리 메모리(리소스 소비 UDF) 또는 특수 리소스 요구사항(GPU 또는 SSD 디스크) 등입니다. (Job 대신) 쿼리마다 힌트로 리소스를 프로파일링하는 것은 매우 유연합니다.
동적 테이블 옵션 (Dynamic Table Options)
동적 테이블 옵션은 테이블 옵션을 동적으로 지정하거나 재정의할 수 있게 합니다. SQL DDL 이나 connect API 로 정의되는 정적 테이블 옵션과 다르게, 이러한 옵션은 각 쿼리 내에서 테이블 단위 범위로 유연하게 지정될 수 있습니다.
따라서 대화형 터미널의 임시(ad-hoc) 쿼리에 사용하기에 매우 적합합니다. 예를 들어 SQL-CLI 에서 동적 옵션 /*+ OPTIONS('csv.ignore-parse-errors'='true') */ 을 추가해 CSV 소스의 파싱 오류를 무시하도록 지정할 수 있습니다.
문법
SQL 호환성을 깨지 않기 위해 Oracle 스타일의 SQL 힌트 문법을 사용합니다:
table_path /*+ OPTIONS(key=val [, key=val]*) */
key:
stringLiteral
val:
stringLiteral
예제
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 Hints 가 지정된 현재 Query block(Query blocks 란 무엇인가?) 입니다. 현재 Flink Query Hints 는 Join Hints 만 지원합니다.
문법
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
쿼리 힌트의 충돌 사례
키-값 힌트 충돌 해결
다음 문법으로 제공되는 키-값(key-value) 힌트의 경우:
hintName '(' optionKey '=' optionVal [, optionKey '=' optionVal ]* ')'
Flink 는 키-값 힌트에서 충돌이 발생하면 마지막 쓰기 승리(last-write-wins) 전략을 채택합니다. 즉 같은 키에 대해 여러 힌트 값이 제공되면 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) 힌트는 다음 문법으로 제공됩니다:
hintName '(' hintOption [, hintOption ]* ')'
목록 힌트로 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 는 사용자가 더 고성능의 실행 계획을 얻기 위해 최적화 도구에 join 전략을 제안할 수 있게 합니다. 현재 Flink Join Hints 는 BROADCAST, SHUFFLE_HASH, SHUFFLE_MERGE, NEST_LOOP 을 지원합니다.
참고:
- Join Hints 에 지정된 테이블은 존재해야 합니다. 그렇지 않으면 테이블 존재하지 않음 오류가 발생합니다.
- 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 전략을 선택합니다. (참고: Join Hints 의 충돌 사례)
BROADCAST
Batch
BROADCAST 는 Flink 가 BroadCast join 을 사용하도록 제안합니다. 힌트가 있는 join 측이 table.optimizer.join.broadcast-threshold 와 무관하게 브로드캐스트됩니다. 따라서 힌트 측 테이블의 데이터 양이 매우 작을 때 잘 수행됩니다.
참고: BROADCAST 는 동치(equivalence) join 조건이 있는 join 만 지원하며 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
Batch
SHUFFLE_HASH 는 Flink 가 Shuffle Hash join 을 사용하도록 제안합니다. 힌트가 있는 join 측이 join 빌드(build) 측이 되며, 힌트 측 테이블의 데이터 양이 너무 크지 않을 때 잘 수행됩니다.
참고: SHUFFLE_HASH 는 동치 join 조건이 있는 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 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
Batch
SHUFFLE_MERGE 는 Flink 가 Sort Merge join 을 사용하도록 제안합니다. 이 유형의 Join Hint 는 두 개의 큰 테이블 사이의 join 시나리오나 join 양측의 데이터가 이미 정렬되어 있는 시나리오에서 사용하는 것이 권장됩니다.
참고: SHUFFLE_MERGE 는 동치 join 조건이 있는 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 (...);
-- 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
Batch
NEST_LOOP 는 Flink 가 Nested Loop join 을 사용하도록 제안합니다. 이 유형의 join 힌트는 특별한 시나리오 요구사항 없이는 권장되지 않습니다.
참고: NEST_LOOP 는 동치 및 비동치 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 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
Streaming
MULTI_JOIN 은 Flink 가 MultiJoin operator 를 사용해 여러 정규 조인을 동시에 처리하도록 제안합니다. 이 유형의 join 힌트는 공통 join 키를 하나 이상 공유하고 큰 중간 상태나 레코드 증폭을 겪는 여러 조인이 있을 때 권장됩니다. MultiJoin operator 는 여러 입력 스트림에 걸쳐 조인을 동시에 처리함으로써 중간 상태를 제거하며, 일부 경우 상태 크기를 크게 줄이고 성능을 개선할 수 있습니다.
MultiJoin operator 에 대한 자세한 내용(사용 시기와 구성 옵션 포함) 은 Multiple Regular Joins 를 참고하세요.
참고:
- MULTI_JOIN 힌트는 테이블 이름 또는 테이블 별칭을 지정할 수 있습니다. 테이블에 별칭이 있으면 힌트는 별칭을 사용해야 합니다.
- MultiJoin operator 가 적용되려면 조인 조건 사이에 키가 하나 이상 공유되어야 합니다.
- 지정되면 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
Streaming
LOOKUP 힌트는 사용자가 Flink 최적화 도구에 다음을 제안할 수 있게 합니다:
- 동기(sync) 또는 비동기(async) 조회 함수 사용
- async 매개변수 구성
- 조회를 위한 지연 재시도 전략 활성화
LOOKUP 힌트 옵션:
| 옵션 유형 | 옵션 이름 | 필수 | 값 유형 | 기본값 | 설명 |
|---|---|---|---|---|---|
| table | table | Y | string | N/A | lookup source 의 테이블 이름 |
| async | async | N | boolean | N/A | 값은 'true' 또는 'false' 일 수 있으며 플래너가 해당 조회 함수를 선택하도록 제안합니다. 백엔드 lookup source 가 제안된 조회 모드를 지원하지 않으면 효과가 없습니다. |
| 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 | lookup join operator 의 백엔드 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 | 사용자 정의 lookup shuffle 을 활성화할지 여부. lookup source 가 입력 데이터 분포를 결정하고 그에 따라 조회 전략을 최적화할 수 있게 합니다. |
참고:
- 'table' 옵션은 필수이며 테이블 이름만 지원됩니다(FROM 절과 일치하게 유지). 테이블에 별칭이 있으면 별칭 이름만 사용할 수 있습니다.
- async 옵션은 모두 선택 사항이며 구성되지 않으면 기본값을 사용합니다.
- retry 옵션에는 기본값이 없습니다. retry 를 활성화해야 할 때 모든 retry 옵션을 유효한 값으로 설정해야 합니다.
1. 동기 및 비동기 조회 함수 사용하기
커넥터에 async 와 sync 조회 기능이 모두 있으면 사용자는 'async'='false' 옵션 값을 주어 플래너가 동기 조회를 사용하도록 제안하거나 'async'='true' 로 비동기 조회를 사용하도록 제안할 수 있습니다:
예제:
-- 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')
참고: async 옵션은 job 수준 실행 옵션 의 async 옵션과 일치하며, 설정하지 않으면 job 수준 구성을 사용합니다. 또 다른 차이점은 LOOKUP 힌트의 범위가 더 작아 현재 lookup 연산에서 힌트 옵션 세트에 해당하는 테이블 이름으로 제한된다는 것입니다(다른 lookup 연산은 LOOKUP 힌트의 영향을 받지 않습니다).
예를 들어 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. 조회를 위한 지연 재시도 전략 활성화
lookup join 의 지연 재시도는 스트림 데이터와 함께 예상치 못한 풍부화(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')
lookup source 에 기능이 하나만 있으면 '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')
이 기능을 최대한 활용하려면 대상 lookup source 가 사용자 정의 shuffle 을 지원해야 합니다. 커넥터 개발자의 경우 LookupTableSource 서브클래스가 SupportsLookupCustomShuffle 를 구현함으로써 이를 달성할 수 있습니다. 소스가 아직 그러한 지원을 제공하지 않더라도 사용자는 이 기능을 먼저 활성화할 수 있으며, Flink 는 최선을 다해 해시 파티셔닝을 적용할 것이며 이 역시 성능 개선을 가져올 것입니다.
추가 참고 사항
재시도에 캐싱 활성화의 영향
FLIP-221 은 lookup source 에 캐싱 지원을 추가하며, PARTIAL 및 FULL 캐싱 모드가 있습니다(NONE 모드는 캐싱을 비활성화함을 의미). FULL 캐싱이 활성화되면 재시도가 전혀 없습니다(lookup source 의 전체 캐시 미러를 통해 조회를 재시도하는 것은 무의미하기 때문). PARTIAL 캐싱이 활성화되면 오는 레코드에 대해 먼저 로컬 캐시에서 조회하고, 캐시 미스 시 백엔드 커넥터를 통해 외부 조회를 수행합니다(캐시 히트 시 레코드를 즉시 반환). 조회 결과가 비어 있으면(캐싱이 비활성화된 것과 동일) 재시도를 트리거하며, 재시도가 완료되면 최종 조회 결과가 결정됩니다(PARTIAL 캐싱 모드에서는 로컬 캐시도 업데이트합니다).
Lookup 키와 'retry-predicate'='lookup_miss' 재시도 조건에 대한 참고 사항
커넥터마다 인덱스 조회 기능이 다를 수 있습니다. 예를 들어 내장 HBase 커넥터는 rowkey 만 조회할 수 있지만(보조 인덱스 없음), 내장 JDBC 커넥터는 임의 열에 대한 더 강력한 인덱스 조회 기능을 제공할 수 있습니다. 이는 서로 다른 물리적 저장소에 의해 결정됩니다.
여기서 언급하는 lookup key 는 인덱스 조회를 위한 필드 또는 필드 조합입니다. lookup join 의 예처럼, c.id 는 조인 조건 "ON o.customer_id = c.id" 의 lookup key 입니다:
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
join 조건을 "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 가 모두 lookup key 로 사용됩니다:
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 만 lookup key 가 될 수 있고, 나머지 join 조건 c.country = 'US' 는 lookup 결과가 반환된 후 평가됩니다:
CREATE TEMPORARY TABLE Customers (
id INT,
name STRING,
country STRING,
zip STRING,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'hbase-2.2',
...
)
따라서 위 쿼리는 'lookup_miss' retry predicate 와 fixed-delay retry strategy 를 활성화할 때 저장소에 따라 다른 재시도 효과를 가집니다.
예를 들어 Customers 테이블에 다음 행이 있다면:
id=100, country='CN'
order 스트림에서 'id=100' 레코드를 처리할 때, 'jdbc' 커넥터에서는 c.id 와 c.country 가 모두 lookup key 로 사용되므로 해당 lookup 결과는 null 입니다(country='CN' 은 조건 c.country = 'US' 를 만족하지 않음). 따라서 재시도가 트리거됩니다.
'hbase' 커넥터에서는 c.id 만 lookup key 로 사용되므로 해당 lookup 결과는 비어 있지 않습니다(레코드 id=100, country='CN' 을 반환). 따라서 재시도가 트리거되지 않습니다(나머지 join 조건 c.country = 'US' 는 반환된 레코드에 대해 참이 아닌 것으로 평가됩니다).
현재 SQL 의미론을 고려해 'lookup_miss' retry predicate 만 제공됩니다. 차원 테이블의 지연 업데이트를 기다려야 할 때(테이블에 이미 과거 버전 레코드가 존재하는 경우가 아니라 없는 경우) 사용자는 두 가지 해결책을 시도할 수 있습니다:
- DataStream Async I/O 의 새 재시도 지원으로 사용자 정의 retry predicate 를 구현합니다(반환된 레코드에 대한 더 복잡한 판단 허용).
- 타임스탬프로 생성된 일종의 데이터 버전을 비교하는 또 다른 join 조건을 추가해 지연 재시도를 활성화합니다. 위 예제에서
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'
)
join 조건에 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)
지연 재시도 조회를 켜면 lookup 노드에서 백프레셔 문제가 발생할 가능성이 더 높으며, 이는 web ui 의 'Task Manager' 페이지의 'Thread Dump' 를 통해 빠르게 확인할 수 있습니다. async 및 sync 조회 각각에서 thread sleep 의 호출 스택이 나타납니다:
- async lookup:
RetryableAsyncLookupFunctionDelegator - sync lookup:
RetryableLookupFunctionDelegator
참고:
- 재시도가 있는 async lookup 은 모든 입력 데이터에 대해 고정 지연 처리를 수행할 수 없습니다(더 가벼운 방법으로 해결해야 함, 예: 소스 소비를 보류하거나 재시도가 있는 sync 조회 사용).
- sync 조회에서 재시도 실행의 지연 대기는 완전히 동기적입니다. 즉 현재 레코드가 완료될 때까지 다음 레코드의 처리가 시작되지 않습니다.
- async 조회에서 'output-mode' 가 'ORDERED' 모드이면 지연 재시도로 인한 백프레셔 확률이 'UNORDERED' 모드보다 높을 수 있습니다. 이 경우 async 'capacity' 를 늘리는 것이 백프레셔를 줄이는 데 효과적이지 않을 수 있으며, 지연 기간을 줄이는 것을 고려해야 할 수 있습니다.
Join Hints 의 충돌 사례
Join Hints 충돌이 발생하면 Flink 는 가장 잘 맞는 것을 선택합니다.
- 먼저
Join Hints는 충돌 해결을 위한 Flink 쿼리 힌트의 로직을 따릅니다(참고: 쿼리 힌트의 충돌 사례) - 하나의 같은 Join Hint 전략 내에서 충돌하면 Flink 는 join 에 대해 첫 번째 일치 테이블을 선택합니다.
- 서로 다른 Join Hints 전략 간의 충돌이면 Flink 는 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 (...);
-- 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 Hints
Streaming
상태 기반 계산 Regular Join 및 Group Aggregation 의 경우 사용자는 STATE_TTL 힌트를 사용해 operator 수준의 Idle State Retention Time 을 지정할 수 있습니다. 이를 통해 위 operator 가 파이프라인 수준 구성인 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은 별칭에 힌트를 지정해야 합니다.- 연쇄(cascade) 조인의 경우 지정된 state TTL 은 첫 번째 조인 operator 의 왼쪽 및 오른쪽 state TTL로, 두 번째 조인 operator 의 오른쪽 state TTL로 해석됩니다(상향식 순서). 두 번째 조인 operator 의 왼쪽 state TTL 은 구성
table.exec.state.ttl에서 가져옵니다. 두 번째 조인 operator 의 왼쪽 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 ...;
- 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 로 간주됩니다.
Query blocks 란?
query block 은 SQL 의 기본 단위입니다. 예를 들어 SQL 문의 인라인 뷰나 하위 쿼리는 외부 쿼리에 대해 별도의 query block 으로 간주됩니다.
예제
SQL 문은 여러 하위 쿼리로 구성될 수 있습니다. 하위 쿼리는 SELECT, INSERT 또는 DELETE 일 수 있습니다. 하위 쿼리는 FROM 절, WHERE 절 또는 UNION 이나 UNION ALL 의 하위 선택에 다른 하위 쿼리를 포함할 수 있습니다.
이러한 서로 다른 하위 쿼리 또는 뷰 유형에 대해 여러 query block 으로 구성될 수 있습니다. 예를 들어:
아래의 간단한 쿼리는 하위 쿼리가 하나뿐이지만 두 개의 query block 이 있습니다. 하나는 외부 SELECT 용이고 다른 하나는 하위 쿼리 SELECT 용입니다.
아래의 쿼리는 union 쿼리로, 두 개의 query block 을 포함합니다. 하나는 첫 번째 SELECT 용이고 다른 하나는 두 번째 SELECT 용입니다.
아래의 쿼리는 뷰를 포함하며, 두 개의 query block 이 있습니다. 하나는 외부 SELECT 용이고 다른 하나는 뷰 용입니다.