연속 쿼리에서의 결정성

연속 쿼리에서의 결정성 (Determinism In Continuous Queries)

이 문서에서 다루는 내용은 다음과 같습니다:

  1. 결정성(Determinism)이란 무엇인가
  2. 모든 배치 처리가 결정적인가
  3. 비결정적 결과를 내는 배치 쿼리 두 가지 예시
  4. 배치 처리에서의 비결정성
  5. 스트리밍 처리에서의 결정성
  6. 스트리밍에서의 비결정성
  7. 스트리밍에서의 비결정적 업데이트
  8. 스트리밍 쿼리에서 비결정적 업데이트의 영향을 제거하는 방법

출처: 문서

본문

1. 결정성이란 무엇인가 (What Is Determinism?)

SQL 표준의 결정성에 대한 설명을 인용하면: "연산이 동일한 입력 값으로 반복될 때 확실히 동일한 결과를 계산하면 그 연산은 결정적(deterministic)이다."

2. 모든 배치 처리가 결정적인가 (Is All Batch Processing Deterministic?)

전형적인 배치 시나리오에서는 주어진 유한 데이터 집합에 대해 동일한 쿼리를 반복 실행하면 일관된 결과가 나옵니다. 이것이 결정성에 대한 가장 직관적인 이해입니다.

그러나 실제로는 배치 처리에서도 동일한 쿼리가 항상 일관된 결과를 반환하지는 않습니다. 두 가지 예시 쿼리를 살펴보겠습니다.

2.1 비결정적 결과를 내는 배치 쿼리 두 가지 예시 (Two Examples Of Batch Queries With Non-Deterministic Results)

예를 들어 새로 만든 웹사이트 클릭 로그 테이블이 있습니다:

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 |
+------+---------------------+------------+
  1. 쿼리 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 |
+------+---------------------+------------+
  1. 쿼리 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 배치 처리에서의 비결정성 (Non-Determinism In Batch Processing)

배치 처리의 비결정성은 주로 비결정적 함수 때문입니다. 위 두 가지 쿼리 예시처럼 내장 함수 CURRENT_TIMESTAMPUUID()는 배치 처리에서 실제로 다르게 동작합니다. 쿼리 예시를 계속 살펴봅니다:

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에는 결정적 함수(deterministic function) 외에도 두 가지 유형의 함수가 있습니다: 비결정적 함수(non-deterministic function)와 동적 함수(dynamic function, 내장 동적 함수는 주로 시간 함수입니다). 비결정적 함수는 런타임에 실행되고(클러스터에서 레코드별로 평가됨), 동적 함수는 쿼리 계획이 생성될 때만 해당 값을 결정합니다(런타임에 실행되지 않으며, 시점에 따라 다른 값을 얻지만 같은 실행에서는 같은 값을 얻습니다). 자세한 내용은 System (Built-in) Function Determinism을 참고하세요.

3. 스트리밍 처리에서의 결정성 (Determinism In Streaming Processing)

스트리밍과 배치의 핵심 차이는 데이터의 무한성(unboundedness)입니다. Flink SQL은 스트리밍 처리를 동적 테이블에 대한 연속 쿼리(continuous query)로 추상화합니다. 따라서 배치 쿼리 예시의 동적 함수는 스트리밍 처리에서 비결정적 함수와 동등합니다(논리적으로 기본 테이블의 모든 변경이 쿼리 실행을 트리거하기 때문입니다). 예시의 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 스트리밍에서의 비결정성 (Non-Determinism In Streaming)

비결정적 함수 외에도 비결정성을 만들 수 있는 다른 요인은 크게 다음과 같습니다:

  1. 소스 커넥터의 비결정적 백 리드(back read)
  2. 처리 시간(Processing Time) 기반 쿼리
  3. TTL 기반 내부 상태 데이터 정리
소스 커넥터의 비결정적 백 리드 (Non-Deterministic Back Read Of Source Connector)

Flink SQL의 경우 제공되는 결정성은 계산에만 한정됩니다. Flink는 사용자 데이터 자체를 저장하지 않기 때문입니다(여기서 스트리밍 모드의 관리형 내부 상태와 사용자 데이터 자체를 구분할 필요가 있습니다). 따라서 결정적 백 리드를 제공할 수 없는 Source 커넥터 구현은 입력 데이터의 비결정성을 가져와 비결정적 결과를 만들 수 있습니다. 흔한 예로 같은 오프셋을 여러 번 읽을 때 데이터가 일관되지 않거나, 보존 시간 때문에 더 이상 존재하지 않는 데이터를 요청하는 경우(예: Kafka 토픽의 구성된 ttl을 넘어선 데이터 요청)가 있습니다.

처리 시간 기반 쿼리 (Query Based On Processing Time)

이벤트 시간과 달리 처리 시간은 머신의 로컬 시간에 기반하며, 이 처리는 결정성을 제공하지 않습니다. 시간 속성에 의존하는 관련 연산으로는 Window Aggregation, Interval Join, Temporal Join 등이 있습니다. 또 다른 전형적인 연산은 Lookup Join으로, 처리 시간 기반 Temporal Join과 의미상 유사하며, 접근하는 외부 테이블이 시간에 따라 변할 때 비결정성이 발생합니다.

TTL 기반 내부 상태 데이터 정리 (Clear Internal State Data Based On TTL)

스트리밍 처리의 무한한 특성 때문에, Regular Join과 Group Aggregation(비윈도우 집계) 같은 연산에서 장기 실행 스트리밍 쿼리가 유지하는 내부 상태 데이터는 계속 커질 수 있습니다. 상태 TTL을 설정해 내부 상태 데이터를 정리하는 것은 종종 필요한 절충이지만, 이것이 계산 결과를 비결정적으로 만들 수도 있습니다.

비결정성이 서로 다른 쿼리에 미치는 영향은 다릅니다. 일부 쿼리에는 단지 비결정적 결과를 만들 뿐이며(쿼리는 정상 동작하지만 여러 번 실행해도 일관된 결과를 내지 못함), 일부 쿼리는 잘못된 결과나 런타임 오류 같은 더 심각한 영향을 받을 수 있습니다. 후자의 주된 이유는 '비결정적 업데이트(non-deterministic update)'입니다.

3.2 스트리밍에서의 비결정적 업데이트 (Non-Deterministic Update In Streaming)

Flink SQL은 '동적 테이블에 대한 연속 쿼리' 추상화에 기반한 완전한 증분 업데이트 메커니즘을 구현합니다. 증분 메시지를 생성해야 하는 모든 연산은 완전한 내부 상태 데이터를 유지하며, 전체 쿼리 파이프라인(소스에서 싱크 연산자까지의 완전한 DAG 포함)의 동작은 연산자 간 업데이트 메시지의 올바른 전달 보장에 의존합니다. 이 보장은 비결정성에 의해 깨질 수 있어 오류를 일으킵니다.

'비결정적 업데이트(Non-deterministic Update, NDU)'란 무엇인가? 업데이트 메시지(changelog)는 여러 종류의 메시지 유형을 포함할 수 있습니다: Insert(I), Delete(D), Update_Before(UB), Update_After(UA). insert-only changelog 파이프라인에는 NDU 문제가 없습니다. 업데이트 메시지(I 외에 D, UB, UA 중 하나 이상을 포함)가 있을 때, 메시지의 업데이트 키(뒤 changelog의 기본 키로 간주할 수 있음)는 쿼리에서 추론됩니다:

  • 업데이트 키를 추론할 수 있으면 파이프라인의 연산자는 업데이트 키로 내부 상태를 유지합니다.
  • 업데이트 키를 추론할 수 없으면(CDC 소스 테이블이나 싱크 테이블에서 기본 키가 정의되지 않았거나, 일부 연산이 쿼리 의미론에서 추론될 수 없는 경우), 내부 상태를 유지하는 모든 연산자는 완전한 행을 통해서만 업데이트(D/UB/UA) 메시지를 처리할 수 있고, 기본 키가 정의되지 않으면 싱크 노드는 retract 모드로 동작하며, 삭제 연산은 완전한 행으로 수행됩니다.

따라서 update-by-row 모드에서 상태를 유지해야 하는 연산자가 받는 모든 업데이트 메시지는 비결정적 컬럼 값에 간섭을 받으면 안 됩니다. 그렇지 않으면 계산 오류를 일으키는 NDU 문제가 발생합니다. 업데이트 메시지가 있고 업데이트 키를 유도할 수 없는 쿼리 파이프라인에서 다음 세 가지가 NDU 문제의 가장 중요한 원인입니다:

  1. 비결정적 함수(스칼라, 테이블, 집계 함수를 포함하고 내장 또는 커스텀 모두)
  2. 변화하는 소스에 대한 LookupJoin
  3. CDC 소스가 메타데이터 필드(시스템 컬럼, 엔티티 행 자체에 속하지 않음)를 운반하는 경우

참고: TTL 기반 내부 상태 데이터 정리로 인한 예외는 별도의 런타임 장애 허용 처리 전략(FLINK-24666)으로 논의됩니다.

3.3 스트리밍에서 비결정적 업데이트의 영향을 제거하는 방법 (How To Eliminate The Impact Of Non-Deterministic Update In Streaming)

스트리밍 쿼리의 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을 조정하라고 안내하는 상세 오류 메시지를 제공합니다(materialization이 가져오는 연산자의 높은 비용과 복잡성을 고려하면 아직 대응하는 자동 해결 메커니즘은 지원되지 않습니다).

모범 사례 (Best Practices)
  1. 스트리밍 쿼리를 실행하기 전에 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 문제를 제거하고 런타임 오류를 피할 수 있습니다.

  1. Lookup Join을 사용할 때는 기본 키를 선언하도록 하세요(존재하는 경우). 기본 키 정의가 있는 Lookup 소스 테이블은 대부분의 경우 Flink SQL이 업데이트 키를 유도하는 것을 막아 비싼 materialization 비용을 절약할 수 있습니다.

다음 두 예시는 사용자가 룩업 소스 테이블에 기본 키를 선언해야 하는 이유를 보여줍니다:

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가 활성화되면 materialization을 추가해 해결할 수 있지만, 비용이 매우 높아 기본 키가 있는 경우보다 두 번 더 비싼 materialization이 생깁니다.

  1. Lookup Join이 접근하는 룩업 소스 테이블이 정적일 때는 TRY_RESOLVE 모드를 활성화하지 않아도 됩니다. Lookup Join이 정적 룩업 소스 테이블에 접근할 때는 먼저 TRY_RESOLVE 모드를 켜서 다른 NDU 문제가 없는지 확인한 뒤, 불필요한 materialization 오버헤드를 피하기 위해 IGNORE 모드로 복원할 수 있습니다.

참고: 여기서 룩업 소스 테이블이 순수하게 정적이고 업데이트되지 않음을 보장해야 합니다. 그렇지 않으면 IGNORE 모드는 안전하지 않습니다.

더 알아보기 (Learn more)