동적 테이블
동적 테이블 (Dynamic Tables)
SQL과 Table API는 실시간 데이터 처리를 위한 유연하고 강력한 기능을 제공합니다. 이 문서는 관계형 개념이 스트리밍에 어떻게 우아하게 적용되는지 설명하며, Flink가 무한 스트림에서 동일한 의미론을 달성할 수 있게 합니다.
출처: 문서
본문
SQL - 그리고 Table API - 은 실시간 데이터 처리를 위한 유연하고 강력한 기능을 제공합니다. 이 문서는 관계형 개념이 스트리밍에 어떻게 우아하게 적용되어 Flink가 무한(unbounded) 스트림에서 동일한 의미론을 달성할 수 있게 하는지 설명합니다.
데이터 스트림에 대한 관계형 질의 (Relational Queries on Data Streams)
다음 표는 입력 데이터, 실행, 출력 결과 측면에서 전통적인 관계형 대수와 스트림 처리를 비교합니다.
| 관계형 대수 / SQL | 스트림 처리 |
|---|---|
| 관계(또는 테이블)는 유한한 튜플의 (멀티)집합입니다. | 스트림은 튜플의 무한 시퀀스입니다. |
| 배치 데이터(예: 관계형 데이터베이스의 테이블)에 실행되는 질의는 완전한 입력 데이터에 접근할 수 있습니다. | 스트리밍 질의는 시작할 때 모든 데이터에 접근할 수 없으며 데이터가 스트리밍되어 들어오기를 "기다려야" 합니다. |
| 배치 질의는 고정된 크기의 결과를 생성한 후 종료됩니다. | 스트리밍 질의는 수신된 레코드에 따라 결과를 지속적으로 갱신하며 결코 완료되지 않습니다. |
이러한 차이에도 불구하고 관계형 질의와 SQL은 스트림 처리를 위한 강력한 도구 모음을 제공합니다. 고급 관계형 데이터베이스 시스템은 Materialized Views라는 기능을 제공합니다. 머티리얼라이즈드 뷰는 일반 가상 뷰처럼 SQL 질의로 정의됩니다. 가상 뷰와 달리 머티리얼라이즈드 뷰는 질의 결과를 캐시하여 접근할 때 질의를 평가할 필요가 없습니다. 캐싱의 일반적인 과제는 캐시가 오래된 결과를 제공하지 않도록 막는 것입니다. 정의 질의의 베이스 테이블이 수정되면 머티리얼라이즈드 뷰는 더 이상 사용할 수 없게 됩니다. Eager View Maintenance는 베이스 테이블이 갱신되는 즉시 머티리얼라이즈드 뷰를 갱신하는 기법입니다.
다음을 고려하면 eager view maintenance와 스트림에 대한 SQL 질의 사이의 연결이 분명해집니다:
- 데이터베이스 테이블은
INSERT,UPDATE,DELETEDML 문의 스트림에서 비롯되며, 흔히 changelog stream이라고 합니다. - 머티리얼라이즈드 뷰는 SQL 질의로 정의됩니다. 뷰를 갱신하려면 질의가 뷰의 베이스 관계의 변경로그 스트림을 지속적으로 처리해야 합니다.
- 머티리얼라이즈드 뷰는 스트리밍 SQL 질의의 결과입니다.
이 점들을 염두에 두고 다음 섹션에서 Dynamic tables 개념을 소개합니다.
동적 테이블과 연속 질의 (Dynamic Tables & Continuous Queries)
동적 테이블(dynamic tables) 은 스트리밍 데이터에 대한 Flink의 Table API 및 SQL 지원의 핵심 개념입니다. 배치 데이터를 나타내는 정적 테이블과 달리 동적 테이블은 시간에 따라 변합니다. 그러나 정적 배치 테이블처럼 시스템은 동적 테이블에 대해 질의를 실행할 수 있습니다. 동적 테이블에 대한 질의는 Continuous Query를 생성합니다. 연속 질의는 종료되지 않으며 동적 결과 - 또 다른 동적 테이블 - 를 생성합니다. 질의는 (동적) 입력 테이블의 변경을 반영하도록 (동적) 결과 테이블을 지속적으로 갱신합니다. 본질적으로 동적 테이블에 대한 연속 질의는 머티리얼라이즈드 뷰를 정의하는 질의와 매우 유사합니다.
중요한 것은 연속 질의의 출력은 항상 입력 테이블의 스냅샷에 대해 배치 모드로 실행한 동일한 질의의 결과와 의미론적으로 동일하다는 것입니다.
다음 그림은 스트림, 동적 테이블, 연속 질의의 관계를 시각화합니다:
- 스트림이 동적 테이블로 변환됩니다.
- 동적 테이블에 연속 질의가 평가되어 새 동적 테이블을 생성합니다.
- 결과 동적 테이블이 다시 스트림으로 변환됩니다.
정보: 동적 테이블은 무엇보다 논리적 개념입니다. 동적 테이블이 질의 실행 중에 반드시 (완전히) 머티리얼라이즈되는 것은 아닙니다.
다음에서 우리는 다음 스키마를 가진 클릭 이벤트 스트림으로 동적 테이블과 연속 질의의 개념을 설명합니다:
CREATE TABLE clicks (
user VARCHAR, -- the name of the user
url VARCHAR, -- the URL that was accessed by the user
cTime TIMESTAMP(3) -- the time when the URL was accessed
) WITH (...);
스트림에 테이블 정의 (Defining a Table on a Stream)
관계형 질의로 스트림을 처리하려면 스트림을 Table로 변환해야 합니다. 개념적으로 스트림의 각 레코드는 결과 테이블에 대한 INSERT 수정으로 해석됩니다. 우리는 INSERT 전용 changelog 스트림에서 테이블을 만들고 있습니다.
다음 그림은 클릭 이벤트 스트림(왼쪽)이 테이블(오른쪽)로 변환되는 방식을 시각화합니다. 클릭 스트림의 더 많은 레코드가 삽입됨에 따라 결과 테이블은 지속적으로 커집니다.
정보: 기억하세요, 스트림에 정의된 테이블은 내부적으로 머티리얼라이즈되지 않습니다.
연속 질의 (Continuous Queries)
연속 질의는 동적 테이블에 대해 평가되어 결과로 새 동적 테이블을 생성합니다. 배치 질의와 달리 연속 질의는 종료되지 않으며 입력 테이블의 갱신에 따라 결과 테이블을 갱신합니다. 어느 시점에서든 연속 질의는 입력 테이블의 스냅샷에 대해 배치 모드로 실행한 동일한 질의의 결과와 의미론적으로 동일합니다.
다음에서 우리는 클릭 이벤트 스트림에 정의된 clicks 테이블에 대한 두 가지 예시 질의를 보여줍니다.
첫 번째 질의는 단순한 GROUP-BY COUNT 집계 질의입니다. clicks 테이블을 user 필드로 그룹화하고 방문한 URL 수를 셉니다. 다음 그림은 clicks 테이블에 추가 행이 갱신됨에 따라 질의가 시간에 따라 어떻게 평가되는지 보여줍니다.
질의가 시작되면 clicks 테이블(왼쪽)은 비어 있습니다. 질의는 첫 번째 행이 삽입될 때 결과 테이블을 계산합니다. 첫 행 [Mary, ./home]이 도착한 후 결과 테이블(오른쪽 상단)은 단일 행 [Mary, 1]로 구성됩니다. 두 번째 행 [Bob, ./cart]가 clicks 테이블에 삽입되면 질의는 결과 테이블을 갱신하고 새 행 [Bob, 1]을 삽입합니다. 세 번째 행 [Mary, ./prod?id=1]은 이미 계산된 결과 행의 갱신을 만들어 [Mary, 1]이 [Mary, 2]로 갱신됩니다. 마지막으로 네 번째 행이 clicks 테이블에 추가되면 질의는 세 번째 행 [Liz, 1]을 결과 테이블에 삽입합니다.
두 번째 질의는 첫 번째와 유사하지만 clicks 테이블을 user 속성 외에도 시간별 텀블링 윈도우로 그룹화한 뒤 URL 수를 셉니다(윈도우 같은 시간 기반 계산은 이후에 논의되는 특별한 시간 속성에 기반합니다). 다시, 그림은 동적 테이블의 변화하는 특성을 시각화하기 위해 서로 다른 시점의 입력과 출력을 보여줍니다.
이전과 마찬가지로 입력 테이블 clicks가 왼쪽에 표시됩니다. 질의는 매시간 연속적으로 결과를 계산하고 결과 테이블을 갱신합니다. clicks 테이블은 12:00:00과 12:59:59 사이의 타임스탬프(cTime)를 가진 네 개의 행을 포함합니다. 질의는 이 입력에서 두 개의 결과 행(각 user마다 하나)을 계산하고 결과 테이블에 추가합니다. 13:00:00과 13:59:59 사이의 다음 윈도우에 대해 clicks 테이블은 세 개의 행을 포함하며, 결과 테이블에 또 다른 두 개의 행이 추가됩니다. 시간이 지나면서 clicks에 더 많은 행이 추가됨에 따라 결과 테이블이 갱신됩니다.
갱신 및 추가 질의 (Update and Append Queries)
두 예시 질의가 꽤 유사해 보이지만(둘 다 그룹화된 카운트 집계를 계산), 한 가지 중요한 측면에서 다릅니다:
- 첫 번째 질의는 이전에 방출된 결과를 갱신합니다. 즉, 결과 테이블을 정의하는 changelog 스트림은
INSERT와UPDATE변경을 포함합니다. - 두 번째 질의는 결과 테이블에만 추가합니다. 즉, 결과 테이블의 changelog 스트림은
INSERT변경만으로 구성됩니다.
질의가 append-only 테이블을 생성하는지 갱신 테이블을 생성하는지에 따라 몇 가지 영향이 있습니다:
- 갱신 변경을 하는 질의는 보통 더 많은 상태를 유지해야 합니다(다음 섹션 참고).
- append-only 테이블을 스트림으로 변환하는 것은 갱신 테이블의 변환과 다릅니다(Table to Stream Conversion 섹션 참고).
질의 제한 (Query Restrictions)
모든 것은 아니지만 많은 의미론적으로 유효한 질의가 스트림에 대한 연속 질의로 평가될 수 있습니다. 일부 질의는 유지해야 하는 상태의 크기 때문이거나 갱신 계산이 너무 비싸서 계산하기에 너무 비쌀 수 있습니다.
- 상태 크기 (State Size): 연속 질의는 무한 스트림에 대해 평가되며 종종 몇 주 또는 몇 달 동안 실행되어야 합니다. 따라서 연속 질의가 처리하는 총 데이터 양은 매우 클 수 있습니다. 이전에 방출된 결과를 갱신해야 하는 질의는 갱신하기 위해 방출된 모든 행을 유지해야 합니다. 예를 들어 첫 번째 예시 질의는 각 사용자의 URL 카운트를 저장하여 카운트를 늘리고 입력 테이블이 새 행을 받을 때 새 결과를 보내야 합니다. 등록된 사용자만 추적한다면 유지해야 할 카운트 수가 너무 높지 않을 수 있습니다. 그러나 등록되지 않은 사용자에게 고유한 사용자 이름이 할당되면 유지할 카운트 수가 시간이 지남에 따라 늘어나 결국 질의가 실패하게 할 수 있습니다.
SELECT user, COUNT(url)
FROM clicks
GROUP BY user;
- 갱신 계산 (Computing Updates): 일부 질의는 단일 입력 레코드가 추가되거나 갱신되어도 방출된 결과 행의 상당 부분을 재계산하고 갱신해야 합니다. 이러한 질의는 연속 질의로 실행하기에 적합하지 않습니다. 예를 들어 다음 질의는 마지막 클릭 시간을 기준으로 각 사용자의
RANK를 계산합니다.clicks테이블이 새 행을 받는 즉시 사용자의lastAction이 갱신되고 새 rank가 계산됩니다. 그러나 두 행이 같은 rank를 가질 수 없으므로 더 낮은 rank의 모든 행도 갱신되어야 합니다.
SELECT user, RANK() OVER (ORDER BY lastAction)
FROM (
SELECT user, MAX(cTime) AS lastAction FROM clicks GROUP BY user
);
Query Configuration 페이지는 연속 질의의 실행을 제어하는 파라미터를 논의합니다. 일부 파라미터는 유지되는 상태의 크기를 결과 정확도와 교환하는 데 사용할 수 있습니다.
테이블-스트림 변환 (Table to Stream Conversion)
동적 테이블은 일반 데이터베이스 테이블처럼 INSERT, UPDATE, DELETE 변경으로 지속적으로 수정될 수 있습니다. 지속적으로 갱신되는 단일 행의 테이블, UPDATE와 DELETE 수정이 없는 insert-only 테이블, 또는 그 사이의 무엇이든 될 수 있습니다.
동적 테이블을 스트림으로 변환하거나 외부 시스템에 쓸 때 이러한 변경은 인코딩되어야 합니다. Flink의 Table API와 SQL은 동적 테이블의 변경을 인코딩하는 세 가지 방식을 지원합니다:
-
Append-only stream:
INSERT변경으로만 수정되는 동적 테이블은 삽입된 행을 방출하여 스트림으로 변환할 수 있습니다. -
Retract stream: 리트랙트(retract) 스트림은 add message와 retract message 두 가지 유형의 메시지를 가진 스트림입니다. 동적 테이블은
INSERT변경을 add message로,DELETE변경을 retract message로,UPDATE변경을 갱신된(이전) 행에 대한 retract message와 갱신하는(새) 행에 대한 추가 메시지로 인코딩하여 리트랙트 스트림으로 변환됩니다. 다음 그림은 동적 테이블을 리트랙트 스트림으로 변환하는 것을 시각화합니다. -
Upsert stream: 업서트(upsert) 스트림은 upsert message와 delete message 두 가지 유형의 메시지를 가진 스트림입니다. 업서트 스트림으로 변환되는 동적 테이블은 (가능한 복합) 고유 키를 요구합니다. 고유 키를 가진 동적 테이블은
INSERT와UPDATE변경을 upsert 메시지로,DELETE변경을 delete 메시지로 인코딩하여 스트림으로 변환됩니다. 스트림을 소비하는 연산자는 메시지를 올바르게 적용하기 위해 고유 키 속성을 인지해야 합니다. 리트랙트 스트림과의 주요 차이는UPDATE변경이 단일 메시지로 인코딩되어 더 효율적이라는 것입니다. 다음 그림은 동적 테이블을 업서트 스트림으로 변환하는 것을 시각화합니다.
동적 테이블을 DataStream으로 변환하는 API는 Common Concepts 페이지에서 논의됩니다. 동적 테이블을 DataStream으로 변환할 때는 append와 retract 스트림만 지원됨을 유의하세요. 동적 테이블을 외부 시스템으로 방출하는 TableSink 인터페이스는 TableSources and TableSinks 페이지에서 논의됩니다.