지속 쿼리에서의 결정성
지속 쿼리에서의 결정성 (Determinism In Continuous Queries)
이 문서는 결정성(determinism)이 무엇인지, 배치 처리와 스트리밍 처리에서 결정성이 어떻게 달라지는지, 그리고 스트리밍 쿼리에서 비결정적 업데이트(NDU) 문제를 어떻게 제거하는지 설명합니다.
출처: 문서
본문
이 문서는 다음에 대해 다룹니다:
- 결정성이란 무엇인가?
- 모든 배치 처리는 결정적인가?
- 비결정적 결과를 내는 배치 쿼리의 두 예
- 배치 처리에서의 비결정성
- 스트리밍 처리에서의 결정성
- 스트리밍에서의 비결정성
- 스트리밍에서의 비결정적 업데이트
- 스트리밍 쿼리에서 비결정적 업데이트의 영향 제거 방법
1. 결정성(Determinism)이란?
SQL 표준의 결정성에 대한 설명을 인용하자면: '동일한 입력 값으로 반복했을 때 확실히 동일한 결과를 계산하는 경우 그 연산은 결정적이다'.
2. 모든 배치 처리는 결정적인가?
전형적인 배치 시나리오에서는 주어진 유한 데이터 집합에 대해 동일한 쿼리를 반복 실행하면 일관된 결과가 나옵니다. 이것이 결정성에 대한 가장 직관적인 이해입니다.
그러나 실제로는 배치 처리에서도 동일한 쿼리가 항상 일관된 결과를 반환하지는 않습니다. 두 가지 예시 쿼리를 살펴보겠습니다.
2.1 비결정적 결과를 내는 배치 쿼리의 두 예
예를 들어 새로 만든 웹사이트 클릭 로그 테이블이 있습니다:
CREATE TABLE clicks (
uid VARCHAR(128),
cTime TIMESTAMP(3),
url VARCHAR(256)
)
일부 새 데이터가 기록되었습니다:
+------+---------------------+------------+
| uid | cTime | url |
+------+---------------------+------------+
| Mary | 2022-08-22 12:00:01 | /home |
| Bob | 2022-08-22 12:00:01 | /home |
| Mary | 2022-08-22 12:00:05 | /prod?id=1 |
| Liz | 2022-08-22 12:01:00 | /home |
| Mary | 2022-08-22 12:01:30 | /cart |
| Bob | 2022-08-22 12:01:35 | /prod?id=3 |
+------+---------------------+------------+
- Query 1 은 로그 테이블에 시간 필터를 적용해 마지막 2분간의 로그를 걸러내려고 합니다:
SELECT * FROM clicks
WHERE cTime BETWEEN TIMESTAMPADD(MINUTE, -2, CURRENT_TIMESTAMP) AND CURRENT_TIMESTAMP;
쿼리가 '2022-08-22 12:02:00' 에 제출되었을 때는 테이블의 6개 행을 모두 반환했고, 1분 후 '2022-08-22 12:03:00' 에 다시 실행했을 때는 3개 항목만 반환했습니다:
+------+---------------------+------------+
| uid | cTime | url |
+------+---------------------+------------+
| Liz | 2022-08-22 12:01:00 | /home |
| Mary | 2022-08-22 12:01:30 | /cart |
| Bob | 2022-08-22 12:01:35 | /prod?id=3 |
+------+---------------------+------------+
- Query 2 는 (
clicks테이블에 기본 키가 없으므로) 반환되는 각 레코드에 고유 식별자를 추가하려고 합니다:
SELECT UUID() AS uuid, * FROM clicks LIMIT 3;
이 쿼리를 연달아 두 번 실행하면 각 행에 대해 다른 uuid 식별자가 생성됩니다:
-- first execution
+--------------------------------+------+---------------------+------------+
| uuid | uid | cTime | url |
+--------------------------------+------+---------------------+------------+
| aaaa4894-16d4-44d0-a763-03f... | Mary | 2022-08-22 12:00:01 | /home |
| ed26fd46-960e-4228-aaf2-0aa... | Bob | 2022-08-22 12:00:01 | /home |
| 1886afc7-dfc6-4b20-a881-b0e... | Mary | 2022-08-22 12:00:05 | /prod?id=1 |
+--------------------------------+------+---------------------+------------+
-- second execution
+--------------------------------+------+---------------------+------------+
| uuid | uid | cTime | url |
+--------------------------------+------+---------------------+------------+
| 95f7301f-bcf2-4b6f-9cf3-1ea... | Mary | 2022-08-22 12:00:01 | /home |
| 63301e2d-d180-4089-876f-683... | Bob | 2022-08-22 12:00:01 | /home |
| f24456d3-e942-43d1-a00f-fdb... | Mary | 2022-08-22 12:00:05 | /prod?id=1 |
+--------------------------------+------+---------------------+------------+
2.2 배치 처리에서의 비결정성
배치 처리에서의 비결정성은 주로 비결정적 함수로 인해 발생합니다. 위의 두 쿼리 예처럼 내장 함수 CURRENT_TIMESTAMP 와 UUID() 는 실제로 배치 처리에서 다르게 동작합니다. 쿼리 예를 계속 살펴보겠습니다:
SELECT CURRENT_TIMESTAMP, * FROM clicks;
CURRENT_TIMESTAMP 는 반환되는 모든 레코드에서 같은 값입니다:
+-------------------------+------+---------------------+------------+
| CURRENT_TIMESTAMP | uid | cTime | url |
+-------------------------+------+---------------------+------------+
| 2022-08-23 17:25:46.831 | Mary | 2022-08-22 12:00:01 | /home |
| 2022-08-23 17:25:46.831 | Bob | 2022-08-22 12:00:01 | /home |
| 2022-08-23 17:25:46.831 | Mary | 2022-08-22 12:00:05 | /prod?id=1 |
| 2022-08-23 17:25:46.831 | Liz | 2022-08-22 12:01:00 | /home |
| 2022-08-23 17:25:46.831 | Mary | 2022-08-22 12:01:30 | /cart |
| 2022-08-23 17:25:46.831 | Bob | 2022-08-22 12:01:35 | /prod?id=3 |
+-------------------------+------+---------------------+------------+
이 차이는 Flink 가 Apache Calcite 로부터 함수의 정의를 상속하기 때문입니다. Calcite 에는 결정적 함수 외에 비결정적 함수(non-deterministic functions)와 동적 함수(dynamic functions, 내장 동적 함수는 주로 시간 함수) 두 가지 유형의 함수가 있습니다. 비결정적 함수는 런타임에 실행되고(클러스터에서는 레코드마다 평가됨), 동적 함수는 쿼리 플랜이 생성될 때만 해당 값을 결정합니다(런타임에 실행되지 않고, 다른 시점에는 다른 값을 얻지만 같은 실행 시점에는 같은 값을 얻습니다). 자세한 내용은 System (Built-in) Function Determinism 을 참고하세요.
3. 스트리밍 처리에서의 결정성
스트리밍과 배치의 핵심 차이는 데이터의 무한성(unboundedness)입니다. Flink SQL 은 스트리밍 처리를 동적 테이블에 대한 지속 쿼리(continuous query on dynamic tables) 로 추상화합니다. 따라서 배치 쿼리 예의 동적 함수는 스트리밍 처리에서 비결정적 함수와 동등합니다(논리적으로 기본 테이블의 모든 변경이 쿼리 실행을 유발합니다). 예의 clicks 로그 테이블이 지속적으로 기록되는 Kafka 토픽에서 온다면, 스트림 모드에서 같은 쿼리는 시간에 따라 변하는 CURRENT_TIMESTAMP 를 반환합니다:
SELECT CURRENT_TIMESTAMP, * FROM clicks;
예를 들어:
+-------------------------+------+---------------------+------------+
| CURRENT_TIMESTAMP | uid | cTime | url |
+-------------------------+------+---------------------+------------+
| 2022-08-23 17:25:46.831 | Mary | 2022-08-22 12:00:01 | /home |
| 2022-08-23 17:25:47.001 | Bob | 2022-08-22 12:00:01 | /home |
| 2022-08-23 17:25:47.310 | Mary | 2022-08-22 12:00:05 | /prod?id=1 |
+-------------------------+------+---------------------+------------+
3.1 스트리밍에서의 비결정성
비결정적 함수 외에도 비결정성을 생성할 수 있는 다른 요인은 주로 다음과 같습니다:
- 소스 커넥터의 비결정적 역판독(back read)
- Processing Time 기반 쿼리
- TTL 기반 내부 상태 데이터 정리
소스 커넥터의 비결정적 역판독
Flink SQL 의 경우 제공되는 결정성은 계산에만 국한됩니다. 자체적으로 사용자 데이터를 저장하지 않기 때문입니다(여기서 스트리밍 모드의 관리되는 내부 상태와 사용자 데이터 자체를 구분할 필요가 있습니다). 따라서 결정적 역판독을 제공할 수 없는 Source 커넥터의 구현은 입력 데이터의 비결정성을 가져와 비결정적 결과를 생성할 수 있습니다. 흔한 예로 같은 오프셋을 여러 번 읽었을 때 데이터가 일치하지 않거나, 보존 시간 때문에 더 이상 존재하지 않는 데이터를 요청하는 경우(예: 요청한 데이터가 Kafka 토픽의 구성된 ttl 을 넘어선 경우)가 있습니다.
Processing Time 기반 쿼리
이벤트 시간과 달리 processing time 은 머신의 로컬 시간을 기반으로 하며, 이 처리는 결정성을 제공하지 않습니다. 시간 속성에 의존하는 관련 연산에는 Window Aggregation, Interval Join, Temporal Join 등이 있습니다. 또 다른 전형적인 연산은 Lookup Join 으로, processing time 기반 Temporal Join과 의미상 유사하며 접근하는 외부 테이블이 시간에 따라 변할 때 비결정성이 발생합니다.
TTL 기반 내부 상태 데이터 정리
스트리밍 처리의 무한성 때문에 Regular Join 및 Group Aggregation (비 윈도 집계) 같은 연산에서 장기 실행 스트리밍 쿼리가 유지하는 내부 상태 데이터가 계속 커질 수 있습니다. 내부 상태 데이터를 정리하기 위해 상태 TTL 을 설정하는 것은 종종 필요한 타협이지만, 이로 인해 계산 결과가 비결정적이 될 수도 있습니다.
비결정성이 서로 다른 쿼리에 미치는 영향은 다릅니다. 일부 쿼리에서는 단순히 비결정적 결과를 생성할 뿐이지만(쿼리는 정상 동작하지만 여러 번 실행해도 일관된 결과를 만들지 못함), 일부 쿼리에서는 잘못된 결과나 런타임 오류 같은 더 심각한 영향을 줄 수 있습니다. 후자의 주된 원인은 '비결정적 업데이트' 입니다.
3.2 스트리밍에서의 비결정적 업데이트 (Non-Deterministic Update)
Flink SQL 은 '동적 테이블에 대한 지속 쿼리' 추상화를 기반으로 완전한 증분 업데이트 메커니즘을 구현합니다. 증분 메시지를 생성해야 하는 모든 연산은 완전한 내부 상태 데이터를 유지하며, 전체 쿼리 파이프라인(source 부터 sink operator 까지의 전체 DAG 포함)의 동작은 operator 간 업데이트 메시지의 올바른 전달 보장에 의존하는데, 이는 비결정성으로 인해 깨질 수 있어 오류로 이어집니다.
'비결정적 업데이트(NDU)' 란 무엇인가요? 업데이트 메시지(changelog) 는 Insert (I), Delete (D), Update_Before (UB), Update_After (UA) 같은 여러 종류의 메시지 유형을 포함할 수 있습니다. insert-only changelog 파이프라인에는 NDU 문제가 없습니다. 업데이트 메시지(I 외에도 D, UB, UA 메시지를 하나 이상 포함)가 있을 때, 메시지의 업데이트 키(업데이트 키는 changelog 의 기본 키로 볼 수 있음)는 쿼리에서 추론됩니다:
- 업데이트 키를 추론할 수 있으면 파이프라인의 operator 는 업데이트 키로 내부 상태를 유지합니다.
- 업데이트 키를 추론할 수 없을 때는(CDC 소스 테이블 또는 Sink 테이블에 기본 키가 정의되지 않았거나, 일부 연산은 쿼리 의미론으로는 추론할 수 없는 경우) 내부 상태를 유지하는 모든 operator 는 완전한 행을 통해서만 업데이트(D/UB/UA) 메시지를 처리할 수 있으며, sink 노드는 기본 키가 정의되지 않으면 retract 모드로 동작하고 삭제 연산은 완전한 행으로 수행됩니다.
따라서 업데이트-행(update-by-row) 모드에서는 상태를 유지해야 하는 operator 가 받는 모든 업데이트 메시지가 비결정적 열 값의 간섭을 받아서는 안 됩니다. 그렇지 않으면 NDU 문제가 발생해 계산 오류가 생길 수 있습니다. 업데이트 메시지가 있고 업데이트 키를 추론할 수 없는 쿼리 파이프라인에서 NDU 문제의 가장 중요한 세 가지 원인은 다음과 같습니다:
- 비결정적 함수(스칼라, 테이블, 집계 함수, 내장형 또는 사용자 정의형 포함)
- 진화하는(evolving) 소스에 대한 Lookup Join
- CDC 소스 가 메타데이터 필드를 운반하는 경우(시스템 열로, 엔티티 행 자체에 속하지 않음)
참고: TTL 기반 내부 상태 데이터 정리로 인한 예외는 별도의 런타임 내결함성 처리 전략으로 논의됩니다 (FLINK-24666).
3.3 스트리밍에서 비결정적 업데이트의 영향 제거 방법
스트리밍 쿼리의 NDU 문제는 보통 직관적이지 않으며, 문제의 위험은 복잡한 쿼리의 작은 변경에서 발생할 수 있습니다.
1.16 부터 Flink SQL (FLINK-27849) 은 실험적 NDU 처리 메커니즘 'table.optimizer.non-deterministic-update.strategy' 을 도입했습니다. TRY_RESOLVE 모드가 활성화되면 스트리밍 쿼리에 NDU 문제가 있는지 확인하고 Lookup Join 이 생성하는 NDU 문제를 제거하려 시도합니다(내부 구체화 materialization 이 추가됨). 위의 1번 또는 3번 지점에 아직 자동으로 제거할 수 없는 요인이 있다면, Flink SQL 은 사용자가 비결정성을 도입하지 않도록 SQL 을 조정하라는 상세한 오류 메시지를 제공합니다(구체화로 인한 operator 의 높은 비용과 복잡성을 고려해, 아직 상응하는 자동 해결 메커니즘은 지원되지 않습니다).
모범 사례 (Best Practices)
- 스트리밍 쿼리를 실행하기 전에
TRY_RESOLVE모드를 활성화하세요. 쿼리에 해결 불가능한 NDU 문제가 있음을 확인하면 오류 프롬프트에 따라 SQL 을 수정해 문제를 적극적으로 피하세요.
FLINK-27639 의 실제 사례:
INSERT INTO t_join_sink
SELECT o.order_id, o.order_name, l.logistics_id, l.logistics_target, l.logistics_source, now()
FROM t_order AS o
LEFT JOIN t_logistics AS l ON ord.order_id=logistics.order_id
기본적으로 생성된 실행 계획은 런타임 오류로 실행됩니다. TRY_RESOLVE 모드가 활성화되면 다음 오류가 제공됩니다:
org.apache.flink.table.api.TableException: The column(s): logistics_time(generated by non-deterministic function: NOW ) can not satisfy the determinism requirement for correctly processing update message(changelogMode contains not only insert 'I').... Please consider removing these non-deterministic columns or making them deterministic.
related rel plan:
Calc(select=[CAST(order_id AS INTEGER) AS order_id, order_name, logistics_id, logistics_target, logistics_source, CAST(NOW() AS TIMESTAMP(6)) AS logistics_time], changelogMode=[I,UB,UA,D], upsertKeys=[])
+- Join(joinType=[LeftOuterJoin], where=[=(order_id, order_id0)], select=[order_id, order_name, logistics_id, logistics_target, logistics_source, order_id0], leftInputSpec=[JoinKeyContainsUniqueKey], rightInputSpec=[HasUniqueKey], changelogMode=[I,UB,UA,D], upsertKeys=[])
:- Exchange(distribution=[hash[order_id]], changelogMode=[I,UB,UA,D], upsertKeys=[{0}])
: +- TableSourceScan(table=[[default_catalog, default_database, t_order, project=[order_id, order_name], metadata=[]]], fields=[order_id, order_name], changelogMode=[I,UB,UA,D], upsertKeys=[{0}])
+- Exchange(distribution=[hash[order_id]], changelogMode=[I,UB,UA,D], upsertKeys=[])
+- TableSourceScan(table=[[default_catalog, default_database, t_logistics, project=[logistics_id, logistics_target, logistics_source, order_id], metadata=[]]], fields=[logistics_id, logistics_target, logistics_source, order_id], changelogMode=[I,UB,UA,D], upsertKeys=[{0}])
오류 프롬프트에 따라 now() 함수를 제거하거나 다른 결정적 함수로 대체하면(또는 order 테이블의 시간 필드를 사용하면) 위의 NDU 문제를 제거하고 런타임 오류를 피할 수 있습니다.
- Lookup Join 을 사용할 때는 (존재한다면) 기본 키를 선언하세요. 기본 키 정의가 있는 Lookup 소스 테이블은 많은 경우 Flink SQL 이 업데이트 키를 도출하는 것을 막아 높은 구체화 비용을 절약할 수 있습니다.
다음 두 예는 왜 Lookup 소스 테이블에 기본 키를 선언하는 것이 권장되는지 보여줍니다:
insert into sink_with_pk
select t1.a, t1.b, t2.c
from (
select *, proctime() proctime from cdc
) t1
join dim_with_pk for system_time as of t1.proctime as t2
on t1.a = t2.a
-- plan: the upsertKey of left stream is reserved when lookup table with a pk definition and use it as lookup key, so that the high cost materialization can be omitted.
Sink(table=[default_catalog.default_database.sink_with_pk], fields=[a, b, c])
+- Calc(select=[a, b, c])
+- LookupJoin(table=[default_catalog.default_database.dim_with_pk], joinType=[InnerJoin], lookup=[a=a], select=[a, b, a, c])
+- DropUpdateBefore
+- TableSourceScan(table=[[default_catalog, default_database, cdc, project=[a, b], metadata=[]]], fields=[a, b])
insert into sink_with_pk
select t1.a, t1.b, t2.c
from (
select *, proctime() proctime from cdc
) t1
join dim_without_pk for system_time as of t1.proctime as t2
on t1.a = t2.a
-- execution plan when `TRY_RESOLVE` is enabled(may encounter errors at runtime when `TRY_RESOLVE` mode is not enabled):
Sink(table=[default_catalog.default_database.sink_with_pk], fields=[a, b, c], upsertMaterialize=[true])
+- Calc(select=[a, b, c])
+- LookupJoin(table=[default_catalog.default_database.dim_without_pk], joinType=[InnerJoin], lookup=[a=a], select=[a, b, a, c], upsertMaterialize=[true])
+- TableSourceScan(table=[[default_catalog, default_database, cdc, project=[a, b], metadata=[]]], fields=[a, b])
두 번째 경우는 TRY_RESOLVE 가 활성화되면 구체화를 추가해 해결할 수 있지만 비용이 매우 높습니다. 기본 키가 있는 경우에 비해 비용이 많이 드는 구체화가 두 번 더 발생합니다.
- Lookup Join 이 접근하는 Lookup 소스 테이블이 정적인 경우
TRY_RESOLVE모드를 활성화하지 않아도 됩니다. Lookup Join 이 정적 Lookup 소스 테이블에 접근할 때는 먼저TRY_RESOLVE모드를 켜서 다른 NDU 문제가 없는지 확인한 다음, 불필요한 구체화 오버헤드를 피하기 위해IGNORE모드로 되돌릴 수 있습니다. 참고: 여기서 Lookup 소스 테이블이 순수하게 정적이고 업데이트되지 않음을 보장해야 합니다. 그렇지 않으면IGNORE모드는 안전하지 않습니다.