CREATE 문
CREATE 문 (CREATE Statements)
CREATE 문은 테이블/뷰/함수를 현재 또는 지정된 Catalog에 등록하는 데 사용돼요. 등록된 테이블/뷰/함수는 SQL 쿼리에서 사용할 수 있어요. Flink SQL이 지원하는 CREATE 문들을 소개해요.
출처: 문서
본문
Flink SQL은 현재 다음 CREATE 문을 지원해요.
- CREATE TABLE
- [CREATE OR] REPLACE TABLE
- CREATE [OR ALTER] MATERIALIZED TABLE
- CREATE CATALOG
- CREATE DATABASE
- CREATE VIEW
- CREATE FUNCTION
- CREATE MODEL
CREATE 문 실행 (Run a CREATE statement)
CREATE 문은 TableEnvironment의 executeSql() 메서드로 실행할 수 있어요. executeSql() 메서드는 성공적인 CREATE 작업이면 'OK'를 반환하고, 그렇지 않으면 예외를 던져요. 다음 예시는 TableEnvironment에서 CREATE 문을 실행하는 방법을 보여줘요. CREATE 문은 SQL CLI에서도 실행할 수 있어요.
TableEnvironment tableEnv = TableEnvironment.create(...);
// SQL query with a registered table
// register a table named "Orders"
tableEnv.executeSql("CREATE TABLE Orders (`user` BIGINT, product STRING, amount INT) WITH (...)");
// run a SQL query on the Table and retrieve the result as a new Table
Table result = tableEnv.sqlQuery(
"SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'");
// Execute insert SQL with a registered table
// register a TableSink
tableEnv.executeSql("CREATE TABLE RubberOrders(product STRING, amount INT) WITH (...)");
// run an insert SQL on the Table and emit the result to the TableSink
tableEnv.executeSql(
"INSERT INTO RubberOrders SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'");
val tableEnv = TableEnvironment.create(...)
// SQL query with a registered table
// register a table named "Orders"
tableEnv.executeSql("CREATE TABLE Orders (`user` BIGINT, product STRING, amount INT) WITH (...)")
// run a SQL query on the Table and retrieve the result as a new Table
val result = tableEnv.sqlQuery(
"SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'")
// Execute insert SQL with a registered table
// register a TableSink
tableEnv.executeSql("CREATE TABLE RubberOrders(product STRING, amount INT) WITH ('connector.path'='/path/to/file' ...)")
// run an insert SQL on the Table and emit the result to the TableSink
tableEnv.executeSql(
"INSERT INTO RubberOrders SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'")
table_env = TableEnvironment.create(...)
# SQL query with a registered table
# register a table named "Orders"
table_env.execute_sql("CREATE TABLE Orders (`user` BIGINT, product STRING, amount INT) WITH (...)");
# run a SQL query on the Table and retrieve the result as a new Table
result = table_env.sql_query(
"SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'");
# Execute an INSERT SQL with a registered table
# register a TableSink
table_env.execute_sql("CREATE TABLE RubberOrders(product STRING, amount INT) WITH (...)")
# run an INSERT SQL on the Table and emit the result to the TableSink
table_env \
.execute_sql("INSERT INTO RubberOrders SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'")
Flink SQL> CREATE TABLE Orders (`user` BIGINT, product STRING, amount INT) WITH (...);
[INFO] Table has been created.
Flink SQL> CREATE TABLE RubberOrders (product STRING, amount INT) WITH (...);
[INFO] Table has been created.
Flink SQL> INSERT INTO RubberOrders SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%';
[INFO] Submitting SQL update statement to the cluster...
CREATE TABLE
다음 문법은 사용 가능한 구문의 개요를 제공해요.
CREATE TABLE [IF NOT EXISTS] [catalog_name.][db_name.]table_name
(
{ <physical_column_definition> | <metadata_column_definition> | <computed_column_definition> }[ , ...n]
[ <watermark_definition> ]
[ <table_constraint> ][ , ...n]
)
[COMMENT table_comment]
[ <distribution> ]
[PARTITIONED BY (partition_column_name1, partition_column_name2, ...)]
WITH (key1=val1, key2=val2, ...)
[ LIKE source_table [( <like_options> )] | AS select_query ]
<physical_column_definition>:
column_name column_type [ <column_constraint> ] [COMMENT column_comment]
<column_constraint>:
[CONSTRAINT constraint_name] PRIMARY KEY NOT ENFORCED
<table_constraint>:
[CONSTRAINT constraint_name] PRIMARY KEY (column_name, ...) NOT ENFORCED
<metadata_column_definition>:
column_name column_type METADATA [ FROM metadata_key ] [ VIRTUAL ]
<computed_column_definition>:
column_name AS computed_column_expression [COMMENT column_comment]
<watermark_definition>:
WATERMARK FOR rowtime_column_name AS watermark_strategy_expression
<source_table>:
[catalog_name.][db_name.]table_name
<like_options>:
{
{ INCLUDING | EXCLUDING } { ALL | CONSTRAINTS | DISTRIBUTION | PARTITIONS }
| { INCLUDING | EXCLUDING | OVERWRITING } { GENERATED | OPTIONS | WATERMARKS }
}[, ...]
<distribution>:
{
DISTRIBUTED BY [ { HASH | RANGE } ] (bucket_column_name1, bucket_column_name2, ...) [INTO n BUCKETS]
| DISTRIBUTED INTO n BUCKETS
}
위 문은 주어진 이름의 테이블을 생성해요. 카탈로그에 같은 이름의 테이블이 이미 존재하면 예외가 발생해요.
컬럼 (Columns)
물리적/일반 컬럼 (Physical / Regular Columns)
물리적 컬럼은 데이터베이스에서 알려진 일반 컬럼이에요. 이들은 물리적 데이터에서 필드의 이름, 유형, 순서를 정의해요. 따라서 물리적 컬럼은 외부 시스템에서 읽고 쓰는 페이로드(payload)를 나타내요. 커넥터와 포맷은 이 컬럼들을 (정의된 순서로) 사용해 스스로를 구성해요. 다른 종류의 컬럼은 물리적 컬럼 사이에 선언될 수 있지만, 최종 물리적 스키마에는 영향을 주지 않아요.
다음 문은 일반 컬럼만 있는 테이블을 생성해요.
CREATE TABLE MyTable (
`user_id` BIGINT,
`name` STRING
) WITH (
...
);
메타데이터 컬럼 (Metadata Columns)
메타데이터 컬럼은 SQL 표준의 확장이며, 테이블의 각 행에 대해 커넥터 및/또는 포맷 특정 필드에 접근할 수 있게 해줘요. 메타데이터 컬럼은 METADATA 키워드로 표시돼요. 예를 들어 메타데이터 컬럼은 시간 기반 연산을 위해 Kafka 레코드의 타임스탬프를 읽고 쓰는 데 사용될 수 있어요. 커넥터와 포맷 문서는 각 구성 요소에 대해 사용 가능한 메타데이터 필드를 나열해요. 그러나 테이블 스키마에서 메타데이터 컬럼을 선언하는 것은 선택 사항이에요.
다음 문은 메타데이터 필드 timestamp를 참조하는 추가 메타데이터 컬럼이 있는 테이블을 생성해요.
CREATE TABLE MyTable (
`user_id` BIGINT,
`name` STRING,
`record_time` TIMESTAMP_LTZ(3) METADATA FROM 'timestamp' -- reads and writes a Kafka record's timestamp
) WITH (
'connector' = 'kafka'
...
);
모든 메타데이터 필드는 문자열 기반 키로 식별되고 문서화된 데이터 유형을 가져요. 예를 들어 Kafka 커넥터는 키가 timestamp이고 데이터 유형이 TIMESTAMP_LTZ(3)인 메타데이터 필드를 노출하며, 이는 레코드 읽기와 쓰기 모두에 사용될 수 있어요.
위 예시에서 메타데이터 컬럼 record_time은 테이블 스키마의 일부가 되고, 일반 컬럼처럼 변환·저장될 수 있어요.
INSERT INTO MyTable SELECT user_id, name, record_time + INTERVAL '1' SECOND FROM MyTable;
편의를 위해, 컬럼 이름이 식별 메타데이터 키로 사용되어야 한다면 FROM 절은 생략할 수 있어요.
CREATE TABLE MyTable (
`user_id` BIGINT,
`name` STRING,
`timestamp` TIMESTAMP_LTZ(3) METADATA -- use column name as metadata key
) WITH (
'connector' = 'kafka'
...
);
편의를 위해, 컬럼의 데이터 유형이 메타데이터 필드의 데이터 유형과 다르면 런타임이 명시적 캐스트를 수행해요. 물론 이는 두 데이터 유형이 호환 가능해야 해요.
CREATE TABLE MyTable (
`user_id` BIGINT,
`name` STRING,
`timestamp` BIGINT METADATA -- cast the timestamp as BIGINT
) WITH (
'connector' = 'kafka'
...
);
기본적으로 플래너는 메타데이터 컬럼이 읽기와 쓰기 모두에 사용될 수 있다고 가정해요. 그러나 대부분의 경우 외부 시스템은 쓰기 가능한 필드보다 더 많은 읽기 전용 메타데이터 필드를 제공해요. 따라서 VIRTUAL 키워드를 사용해 메타데이터 컬럼을 지속(persisting)에서 제외할 수 있어요.
CREATE TABLE MyTable (
`timestamp` BIGINT METADATA, -- part of the query-to-sink schema
`offset` BIGINT METADATA VIRTUAL, -- not part of the query-to-sink schema
`user_id` BIGINT,
`name` STRING,
) WITH (
'connector' = 'kafka'
...
);
위 예시에서 offset은 읽기 전용 메타데이터 컬럼이며 query-to-sink 스키마에서 제외돼요. 따라서 source-to-query 스키마(SELECT용)와 query-to-sink(INSERT INTO용) 스키마가 달라요.
source-to-query schema:
MyTable(`timestamp` BIGINT, `offset` BIGINT, `user_id` BIGINT, `name` STRING)
query-to-sink schema:
MyTable(`timestamp` BIGINT, `user_id` BIGINT, `name` STRING)
계산 컬럼 (Computed Columns)
계산 컬럼은 column_name AS computed_column_expression 구문을 사용해 생성되는 가상 컬럼이에요. 계산 컬럼은 같은 테이블에 선언된 다른 컬럼을 참조할 수 있는 표현식을 평가해요. 물리적 컬럼과 메타데이터 컬럼 모두 접근할 수 있어요. 컬럼 자체는 물리적으로 테이블에 저장되지 않아요. 컬럼의 데이터 유형은 주어진 표현식에서 자동으로 파생되며 수동으로 선언할 필요가 없어요.
플래너는 계산 컬럼을 소스 이후의 일반 프로젝션으로 변환해요. 최적화 또는 워터마크 전략 푸시다운을 위해, 평가는 연산자에 걸쳐 분산되거나, 여러 번 수행되거나, 주어진 쿼리에 필요하지 않으면 건너뛸 수 있어요.
예를 들어 계산 컬럼은 다음과 같이 정의될 수 있어요.
CREATE TABLE MyTable (
`user_id` BIGINT,
`price` DOUBLE,
`quantity` DOUBLE,
`cost` AS price * quantity -- evaluate expression and supply the result to queries
) WITH (
'connector' = 'kafka'
...
);
표현식은 컬럼, 상수, 또는 함수의 어떤 조합이든 포함할 수 있어요. 표현식은 하위 쿼리를 포함할 수 없어요.
계산 컬럼은 Flink에서 CREATE TABLE 문의 **시간 속성(time attributes)**을 정의하는 데 흔히 사용돼요.
처리 시간(processing time) 속성은 시스템의 PROCTIME() 함수를 사용해 proc AS PROCTIME()으로 쉽게 정의할 수 있어요. 이벤트 시간(事件 time) 속성 timestamp는 WATERMARK 선언 전에 사전 처리될 수 있어요. 예를 들어 원본 필드가 TIMESTAMP(3) 유형이 아니거나 JSON 문자열에 중첩되어 있으면 계산 컬럼을 사용할 수 있어요.
가상 메타데이터 컬럼과 유사하게, 계산 컬럼은 지속에서 제외돼요. 따라서 계산 컬럼은 INSERT INTO 문의 대상이 될 수 없어요. 따라서 source-to-query 스키마(SELECT용)와 query-to-sink(INSERT INTO용) 스키마가 달라요.
source-to-query schema:
MyTable(`user_id` BIGINT, `price` DOUBLE, `quantity` DOUBLE, `cost` DOUBLE)
query-to-sink schema:
MyTable(`user_id` BIGINT, `price` DOUBLE, `quantity` DOUBLE)
WATERMARK
WATERMARK 절은 테이블의 이벤트 시간 속성을 정의하며 WATERMARK FOR rowtime_column_name AS watermark_strategy_expression의 형태를 가져요.
rowtime_column_name은 테이블의 이벤트 시간 속성으로 표시되는 기존 컬럼을 정의해요. 컬럼은 TIMESTAMP(3) 유형이어야 하고 스키마의 최상위 컬럼이어야 해요. 계산 컬럼일 수 있어요.
watermark_strategy_expression은 워터마크 생성 전략을 정의해요. 워터마크를 계산하기 위해 계산 컬럼을 포함한 임의의 비-쿼리 표현식을 허용해요. 표현식 반환 유형은 TIMESTAMP(3)이어야 하며, 이는 Epoch 이후의 타임스탬프를 나타내요.
반환되는 워터마크는 null이 아니고 그 값이 이전에 방출된 로컬 워터마크보다 클 때만 방출돼요 (오름차순 워터마크의 계약을 보존하기 위해). 워터마크 생성 표현식은 프레임워크가 모든 레코드에 대해 평가해요.
프레임워크는 주기적으로 가장 큰 생성 워터마크를 방출해요. 현재 워터마크가 이전 것과 동일하거나, null이거나, 반환된 워터마크의 값이 마지막으로 방출된 것보다 작으면 새 워터마크는 방출되지 않아요.
워터마크는 pipeline.auto-watermark-interval 구성으로 정의된 간격으로 방출돼요. 워터마크 간격이 0ms이면, 생성된 워터마크는 null이 아니고 마지막 방출된 것보다 큰 경우 레코드마다 방출돼요.
이벤트 시간 의미론을 사용할 때 테이블은 이벤트 시간 속성과 워터마킹 전략을 포함해야 해요. Flink는 몇 가지 흔히 사용되는 워터마크 전략을 제공해요.
- Strictly ascending timestamps (엄격히 오름차순 타임스탬프):
WATERMARK FOR rowtime_column AS rowtime_column. 지금까지 관찰된 최대 타임스탬프의 워터마크를 방출해요. 최대 타임스탬프보다 큰 타임스탬프를 가진 행은 늦지 않아요. - Ascending timestamps (오름차순 타임스탬프):
WATERMARK FOR rowtime_column AS rowtime_column - INTERVAL '0.001' SECOND. 지금까지 관찰된 최대 타임스탬프에서 1을 뺀 워터마크를 방출해요. 최대 타임스탬프와 같거나 큰 타임스탬프를 가진 행은 늦지 않아요. - Bounded out of orderness timestamps (유한 무순서 타임스탬프):
WATERMARK FOR rowtime_column AS rowtime_column - INTERVAL 'string' timeUnit. 최대 관찰 타임스탬프에서 지정된 지연을 뺀 워터마크를 방출해요. 예를 들어WATERMARK FOR rowtime_column AS rowtime_column - INTERVAL '5' SECOND는 5초 지연 워터마크 전략이에요.
CREATE TABLE Orders (
`user` BIGINT,
product STRING,
order_time TIMESTAMP(3),
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH ( . . . );
PRIMARY KEY
기본 키 제약은 Flink가 최적화에 활용하도록 하는 힌트예요. 테이블 또는 뷰의 컬럼 또는 컬럼 집합이 고유하고 null을 포함하지 않음을 알려줘요. 기본 키의 어떤 컬럼도 nullable일 수 없어요. 따라서 기본 키는 테이블의 행을 고유하게 식별해요.
기본 키 제약은 컬럼 정의와 함께(컬럼 제약) 또는 한 줄로(테이블 제약) 선언될 수 있어요. 두 경우 모두 단일 행으로만 선언되어야 해요. 동시에 여러 기본 키 제약을 정의하면 예외가 발생해요.
유효성 검사 (Validity Check)
SQL 표준은 제약이 ENFORCED 또는 NOT ENFORCED일 수 있다고 명시해요. 이는 들어오는/나가는 데이터에 대해 제약 검사가 수행되는지 여부를 제어해요. Flink는 데이터를 소유하지 않으므로 지원하려는 유일한 모드는 NOT ENFORCED 모드예요. 쿼리가 키 무결성을 강제하도록 보장하는 것은 사용자의 몫이에요.
Flink는 컬럼의 nullable 상태가 기본 키의 컬럼과 정렬되어 있다고 가정해 기본 키의 정확성을 가정해요. 커넥터는 이들이 정렬되도록 보장해야 해요.
참고: CREATE TABLE 문에서 기본 키 제약을 만들면 컬럼의 nullable 상태가 변경돼요. 즉 기본 키 제약이 있는 컬럼은 nullable이 아니게 돼요.
PARTITIONED BY
생성된 테이블을 지정된 컬럼으로 파티셔닝해요. 이 테이블이 파일시스템 싱크로 사용되면 각 파티션에 대해 디렉터리가 생성돼요.
DISTRIBUTED
버킷(Buckets)은 데이터를 분리된 부분 집합으로 나눠 외부 저장 시스템에서 로드 밸런싱을 가능하게 해요. 이러한 부분 집합은 잠재적으로 "무한"한 키스페이스를 가진 행들을 더 작고 관리하기 쉬운 청크로 그룹화해 효율적인 병렬 처리를 허용해요.
버케팅은 기본 커넥터의 의미론에 크게 의존해요. 하지만 사용자는 버킷 수, 버케팅 알고리즘, 그리고 (알고리즘이 허용한다면) 대상 버킷 계산에 사용되는 컬럼을 지정해 버케팅 동작에 영향을 줄 수 있어요.
모든 버케팅 구성 요소(즉 버킷 수, 분배 알고리즘, 버킷 키 컬럼)는 SQL 구문 관점에서 선택 사항이에요.
다음 SQL 문이 주어졌을 때:
-- Example 1
CREATE TABLE MyTable (uid BIGINT, name STRING) DISTRIBUTED BY HASH(uid) INTO 4 BUCKETS;
-- Example 2
CREATE TABLE MyTable (uid BIGINT, name STRING) DISTRIBUTED BY (uid) INTO 4 BUCKETS;
-- Example 3
CREATE TABLE MyTable (uid BIGINT, name STRING) DISTRIBUTED BY (uid);
-- Example 4
CREATE TABLE MyTable (uid BIGINT, name STRING) DISTRIBUTED INTO 4 BUCKETS;
예시 1은 고정된 4개의 버킷에 해시 함수를 선언해요 (즉 HASH(uid) % 4 = target bucket). 예시 2는 알고리즘 선택을 커넥터에 맡겨요. 게다가 예시 3은 버킷 수를 커넥터에 맡겨요. 반면 예시 4는 버킷 수만 정의해요.
WITH 옵션 (WITH Options)
테이블 소스/싱크를 생성하는 데 사용되는 테이블 속성이에요. 이 속성들은 보통 기본 커넥터를 찾고 생성하는 데 사용돼요. key1=val1 표현식의 key와 value 모두 문자열 리터럴이어야 해요. 서로 다른 커넥터의 모든 지원 테이블 속성에 대한 자세한 내용은 Connect to External Systems를 참고하세요.
참고: 테이블 이름은 세 가지 형식일 수 있어요: 1. catalog_name.db_name.table_name 2. db_name.table_name 3. table_name. catalog_name.db_name.table_name의 경우 테이블은 "catalog_name"이라는 이름의 카탈로그와 "db_name"이라는 이름의 데이터베이스로 metastore에 등록돼요. db_name.table_name의 경우 실행 테이블 환경의 현재 카탈로그와 "db_name"이라는 데이터베이스로 등록돼요. table_name의 경우 실행 테이블 환경의 현재 카탈로그와 데이터베이스로 등록돼요.
참고: CREATE TABLE 문으로 등록된 테이블은 테이블 소스와 테이블 싱크 모두로 사용될 수 있어요. DML에서 참조될 때까지 소스로 사용되는지 싱크로 사용되는지 결정할 수 없어요.
LIKE
LIKE 절은 SQL 기능(Feature T171 "테이블 정의의 LIKE 절" 및 Feature T173 "테이블 정의의 확장 LIKE 절")의 변형/결합이에요. 이 절은 기존 테이블의 정의를 기반으로 테이블을 생성하는 데 사용될 수 있어요. 게다가 사용자는 원본 테이블을 확장하거나 특정 부분을 제외할 수 있어요. SQL 표준과 달리 이 절은 CREATE 문의 최상위에 정의되어야 해요. 이는 이 절이 정의의 여러 부분이 아닌(스키마 부분뿐만 아니라) 적용되기 때문이에요.
이 절을 사용해 특정 커넥터 속성을 재사용(그리고 잠재적으로 덮어쓰기)하거나 외부에서 정의된 테이블에 워터마크를 추가할 수 있어요. 예를 들어 Apache Hive에 정의된 테이블에 워터마크를 추가할 수 있어요.
아래 예시 문을 고려해 봅시다.
CREATE TABLE Orders (
`user` BIGINT,
product STRING,
order_time TIMESTAMP(3)
) WITH (
'connector' = 'kafka',
'scan.startup.mode' = 'earliest-offset'
);
CREATE TABLE Orders_with_watermark (
-- Add watermark definition
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
-- Overwrite the startup-mode
'scan.startup.mode' = 'latest-offset'
)
LIKE Orders;
결과 테이블 Orders_with_watermark는 다음 문으로 생성된 테이블과 동일해요.
CREATE TABLE Orders_with_watermark (
`user` BIGINT,
product STRING,
order_time TIMESTAMP(3),
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'scan.startup.mode' = 'latest-offset'
);
테이블 기능의 병합 논리는 like 옵션으로 제어할 수 있어요. 다음의 병합 동작을 제어할 수 있어요.
CONSTRAINTS- 기본 키와 고유 키 같은 제약GENERATED- 계산 컬럼METADATA- 메타데이터 컬럼OPTIONS- 커넥터 및 포맷 속성을 설명하는 커넥터 옵션DISTRIBUTION- 분배 정의PARTITIONS- 테이블의 파티션WATERMARKS- 워터마크 선언
세 가지 병합 전략으로:
INCLUDING- 소스 테이블의 기능을 포함하고, 중복 항목(예: 두 테이블에 같은 키의 옵션이 있는 경우)에서 실패해요.EXCLUDING- 소스 테이블의 주어진 기능을 포함하지 않아요.OVERWRITING- 소스 테이블의 기능을 포함하고, 소스 테이블의 중복 항목을 새 테이블의 속성으로 덮어써요. 예: 두 테이블에 같은 키의 옵션이 있으면 현재 문의 것이 사용돼요.
또한 INCLUDING/EXCLUDING ALL 옵션을 사용해 특정 전략이 정의되지 않은 경우의 전략을 지정할 수 있어요. 즉 EXCLUDING ALL INCLUDING WATERMARKS를 사용하면 소스 테이블에서 워터마크만 포함돼요.
예시:
-- A source table stored in a filesystem
CREATE TABLE Orders_in_file (
`user` BIGINT,
product STRING,
order_time_string STRING,
order_time AS to_timestamp(order_time_string)
)
PARTITIONED BY (`user`)
WITH (
'connector' = 'filesystem',
'path' = '...'
);
-- A corresponding table we want to store in kafka
CREATE TABLE Orders_in_kafka (
-- Add watermark definition
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
...
)
LIKE Orders_in_file (
-- Exclude everything besides the computed columns which we need to generate the watermark for.
-- We do not want to have the partitions or filesystem options as those do not apply to kafka.
EXCLUDING ALL
INCLUDING GENERATED
);
like 옵션을 제공하지 않으면 INCLUDING ALL OVERWRITING OPTIONS가 기본값으로 사용돼요. 참고: 물리적 컬럼의 병합 동작은 제어할 수 없어요. 이들은 INCLUDING 전략을 적용한 것처럼 병합돼요. 참고: source_table은 복합 식별자일 수 있어요. 따라서 다른 카탈로그나 데이터베이스의 테이블일 수 있어요. 예: my_catalog.my_db.MyTable은 카탈로그 my_catalog와 데이터베이스 my_db의 테이블 MyTable을 지정하고, my_db.MyTable은 현재 카탈로그와 데이터베이스 my_db의 테이블 MyTable을 지정해요.
AS select_statement (CTAS)
테이블은 하나의 create-table-as-select(CTAS) 문에서 쿼리 결과로 생성되고 채워질 수도 있어요. CTAS는 단일 명령으로 테이블을 만들고 데이터를 삽입하는 가장 간단하고 빠른 방법이에요.
CTAS에는 두 부분이 있으며, SELECT 부분은 Flink SQL이 지원하는 모든 SELECT 쿼리일 수 있어요. CREATE 부분은 SELECT 부분에서 결과 스키마를 가져와 대상 테이블을 생성해요. CREATE TABLE과 유사하게, CTAS는 대상 테이블의 필수 옵션이 WITH 절에 지정되어야 함을 요구해요.
CTAS의 테이블 생성 연산은 대상 Catalog에 의존해요. 예를 들어 Hive Catalog는 물리적 테이블을 Hive에 자동으로 생성해요. 그러나 인메모리 카탈로그는 SQL이 실행되는 클라이언트의 메모리에 테이블 메타데이터를 등록해요.
아래 예시 문을 고려해 봅시다.
CREATE TABLE my_ctas_table
WITH (
'connector' = 'kafka',
...
)
AS SELECT id, name, age FROM source_table WHERE mod(id, 10) = 0;
결과 테이블 my_ctas_table은 다음 문으로 테이블을 생성하고 데이터를 삽입한 것과 동일해요.
CREATE TABLE my_ctas_table (
id BIGINT,
name STRING,
age INT
) WITH (
'connector' = 'kafka',
...
);
INSERT INTO my_ctas_table SELECT id, name, age FROM source_table WHERE mod(id, 10) = 0;
CREATE 부분은 명시적 컬럼을 지정할 수 있게 해요. 결과 테이블 스키마는 CREATE 부분에 정의된 컬럼을 먼저 포함하고 그다음 SELECT 부분의 컬럼을 포함해요. 두 부분(CREATE와 SELECT) 모두에 이름이 있는 컬럼은 SELECT 부분에 정의된 것과 같은 컬럼 위치를 유지해요. SELECT 컬럼의 데이터 유형도 CREATE 부분에 지정하면 덮어쓸 수 있어요.
아래 예시 문을 고려해 봅시다.
CREATE TABLE my_ctas_table (
desc STRING,
quantity DOUBLE,
cost AS price * quantity,
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND,
) WITH (
'connector' = 'kafka',
...
) AS SELECT id, price, quantity, order_time FROM source_table;
결과 테이블 my_ctas_table은 다음 테이블을 생성하고 다음 문으로 데이터를 삽입한 것과 동일해요.
CREATE TABLE my_ctas_table (
desc STRING,
cost AS price * quantity,
id BIGINT,
price DOUBLE,
quantity DOUBLE,
order_time TIMESTAMP(3),
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
...
);
INSERT INTO my_ctas_table (id, price, quantity, order_time)
SELECT id, price, quantity, order_time FROM source_table;
CREATE 부분은 기본 키와 분배 전략도 지정할 수 있게 해요. 기본 키는 NOT NULL 컬럼에서만 동작한다는 점에 유의하세요. 현재 기본 키는 SELECT 부분의 NOT NULL일 수 있는 컬럼만 정의할 수 있게 해요. CREATE 부분은 NOT NULL 컬럼 정의를 허용하지 않아요.
SELECT 부분에서 id가 not null 컬럼인 아래 예시 문을 고려해 봅시다.
CREATE TABLE my_ctas_table (
PRIMARY KEY (id) NOT ENFORCED
) DISTRIBUTED BY (id) INTO 4 buckets
AS SELECT id, name FROM source_table;
결과 테이블 my_ctas_table은 다음 테이블을 생성하고 다음 문으로 데이터를 삽입한 것과 동일해요.
CREATE TABLE my_ctas_table (
id BIGINT NOT NULL PRIMARY KEY NOT ENFORCED,
name STRING
) DISTRIBUTED BY (id) INTO 4 buckets;
INSERT INTO my_ctas_table SELECT id, name FROM source_table;
CTAS는 또한 CREATE 부분에서 데이터 유형 없이 모든 컬럼 이름을 지정해 SELECT 부분에 정의된 컬럼의 순서를 재정렬할 수 있게 해요. 이 기능은 INSERT INTO 문과 동일해요. 지정된 컬럼은 SELECT 부분의 컬럼 이름과 수와 일치해야 해요. 이 정의는 데이터 유형 정의가 필요하므로 새 컬럼과 결합할 수 없어요.
아래 예시 문을 고려해 봅시다.
CREATE TABLE my_ctas_table (
order_time, price, quantity, id
) WITH (
'connector' = 'kafka',
...
) AS SELECT id, price, quantity, order_time FROM source_table;
결과 테이블 my_ctas_table은 다음 테이블을 생성하고 다음 문으로 데이터를 삽입한 것과 동일해요.
CREATE TABLE my_ctas_table (
order_time TIMESTAMP(3),
price DOUBLE,
quantity DOUBLE,
id BIGINT
) WITH (
'connector' = 'kafka',
...
);
INSERT INTO my_ctas_table (order_time, price, quantity, id)
SELECT id, price, quantity, order_time FROM source_table;
참고: CTAS에는 다음 제한이 있어요.
- 아직 임시 테이블 생성을 지원하지 않아요.
- 아직 파티션 테이블 생성을 지원하지 않아요.
참고: 기본적으로 CTAS는 비원자적(non-atomic)이에요. 즉 테이블에 데이터를 삽입하는 동안 오류가 발생해도 생성된 테이블이 자동으로 삭제되지 않아요.
원자성 (Atomicity)
CTAS에 원자성을 활성화하려면 다음을 확인해야 해요.
- 싱크가 CTAS에 대한 원자성 의미론을 구현했어야 해요. 해당 커넥터 싱크에 대한 문서를 참고해 원자성 의미론이 가능한지 알 수 있어요. 원자성 의미론을 구현하려는 개발자는 SupportsStaging 문서를 참고하세요.
table.rtas-ctas.atomicity-enabled옵션을true로 설정해요.
[CREATE OR] REPLACE TABLE (RTAS)
[CREATE OR] REPLACE TABLE [catalog_name.][db_name.]table_name
[(
{ <physical_column_definition> | <metadata_column_definition> | <computed_column_definition> }[ , ...n]
[ <watermark_definition> ]
[ <table_constraint> ][ , ...n]
)]
[COMMENT table_comment]
[ <distribution> ]
WITH (key1=val1, key2=val2, ...)
AS select_query
참고: RTAS에는 다음 의미론이 있어요.
REPLACE TABLE AS SELECT문: 교체될 대상 테이블이 반드시 존재해야 해요. 그렇지 않으면 예외가 발생해요.CREATE OR REPLACE TABLE AS SELECT문: 교체될 대상 테이블이 존재하지 않으면 생성되고, 존재하면 교체돼요.
테이블은 하나의 [CREATE OR] REPLACE TABLE AS SELECT(RTAS) 문에서 쿼리 결과로 교체(또는 생성)되고 채워질 수 있어요. RTAS는 단일 명령으로 테이블을 교체(또는 생성)하고 데이터를 삽입하는 가장 간단하고 빠른 방법이에요.
RTAS에는 두 부분이 있으며, SELECT 부분은 Flink SQL이 지원하는 모든 SELECT 쿼리일 수 있고, REPLACE TABLE 부분은 SELECT 부분에서 결과 스키마를 가져와 대상 테이블을 교체해요. CREATE TABLE 및 CTAS와 유사하게, RTAS는 대상 테이블의 필수 옵션이 WITH 절에 지정되어야 함을 요구해요.
아래 예시 문을 고려해 봅시다.
REPLACE TABLE my_rtas_table
WITH (
'connector' = 'kafka',
...
)
AS SELECT id, name, age FROM source_table WHERE mod(id, 10) = 0;
REPLACE TABLE AS SELECT 문은 먼저 테이블을 삭제한 다음 다음 문으로 테이블을 생성하고 데이터를 삽입하는 것과 동일해요.
DROP TABLE my_rtas_table;
CREATE TABLE my_rtas_table (
id BIGINT,
name STRING,
age INT
) WITH (
'connector' = 'kafka',
...
);
INSERT INTO my_rtas_table SELECT id, name, age FROM source_table WHERE mod(id, 10) = 0;
CREATE TABLE AS와 유사하게, REPLACE TABLE AS는 명시적 컬럼, 워터마크, 기본 키, 분배 전략을 지정할 수 있게 해요. 결과 테이블 스키마는 CREATE 부분에서 먼저 만들어진 다음 SELECT 부분의 컬럼이 이어져요. 두 부분(CREATE와 SELECT) 모두에 이름이 있는 컬럼은 SELECT 부분에 정의된 것과 같은 컬럼 위치를 유지해요. SELECT 컬럼의 데이터 유형도 CREATE 부분에 지정하면 덮어쓸 수 있어요.
아래 예시 문을 고려해 봅시다.
REPLACE TABLE my_rtas_table (
desc STRING,
quantity DOUBLE,
cost AS price * quantity,
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND,
PRIMARY KEY (id) NOT ENFORCED
) DISTRIBUTED BY (id) INTO 4 buckets
AS SELECT id, price, quantity, order_time FROM source_table;
결과 테이블 my_rtas_table은 다음 문으로 테이블을 생성하고 데이터를 삽입한 것과 동일해요.
DROP TABLE my_rtas_table;
CREATE TABLE my_rtas_table (
desc STRING,
cost AS price * quantity,
id BIGINT NOT NULL PRIMARY KEY NOT ENFORCED,
price DOUBLE,
quantity DOUBLE,
order_time TIMESTAMP(3),
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
...
);
INSERT INTO my_rtas_table (id, price, quantity, order_time)
SELECT id, price, quantity, order_time FROM source_table;
참고: RTAS에는 다음 제한이 있어요.
- 아직 임시 테이블 교체를 지원하지 않아요.
- 아직 파티션 테이블 생성을 지원하지 않아요.
참고: 기본적으로 RTAS는 비원자적(non-atomic)이에요. 즉 테이블에 데이터를 삽입하는 동안 오류가 발생해도 테이블이 자동으로 삭제되거나 원래 상태로 복원되지 않아요. 참고: RTAS는 먼저 테이블을 삭제한 다음 테이블을 생성하고 데이터를 삽입해요. 하지만 테이블이 인메모리 카탈로그에 있으면, 테이블 삭제는 물리적 테이블의 데이터를 제거하지 않고 카탈로그에서만 제거해요. 따라서 RTAS 문 실행 전의 데이터는 여전히 존재해요.
원자성 (Atomicity)
RTAS에 원자성을 활성화하려면 다음을 확인해야 해요.
- 싱크가 RTAS에 대한 원자성 의미론을 구현했어야 해요. 해당 커넥터 싱크에 대한 문서를 참고해 원자성 의미론이 가능한지 알 수 있어요. 원자성 의미론을 구현하려는 개발자는 SupportsStaging 문서를 참고하세요.
table.rtas-ctas.atomicity-enabled옵션을true로 설정해요.
CREATE [OR ALTER] MATERIALIZED TABLE
Materialized 테이블에 대한 전용 페이지를 참고하세요.
CREATE CATALOG
CREATE CATALOG [IF NOT EXISTS] catalog_name
[COMMENT catalog_comment]
WITH (key1=val1, key2=val2, ...)
주어진 카탈로그 속성으로 카탈로그를 생성해요. 같은 이름의 카탈로그가 이미 존재하면 예외가 발생해요.
IF NOT EXISTS: 카탈로그가 이미 존재하면 아무 일도 일어나지 않아요.WITH OPTIONS: 이 카탈로그와 관련된 추가 정보를 저장하는 데 사용되는 카탈로그 속성이에요.key1=val1표현식의 key와 value 모두 문자열 리터럴이어야 해요. 자세한 내용은 Catalogs를 참고하세요.
CREATE DATABASE
CREATE DATABASE [IF NOT EXISTS] [catalog_name.]db_name
[COMMENT database_comment]
[WITH (key1=val1, key2=val2, ...)]
주어진 데이터베이스 속성으로 데이터베이스를 생성해요. 카탈로그에 같은 이름의 데이터베이스가 이미 존재하면 예외가 발생해요.
IF NOT EXISTS: 데이터베이스가 이미 존재하면 아무 일도 일어나지 않아요.WITH OPTIONS: 이 데이터베이스와 관련된 추가 정보를 저장하는 데 사용되는 데이터베이스 속성이에요.key1=val1표현식의 key와 value 모두 문자열 리터럴이어야 해요.
CREATE VIEW
CREATE [TEMPORARY] VIEW [IF NOT EXISTS] [catalog_name.][db_name.]view_name
[( columnName [, columnName ]* )] [COMMENT view_comment]
AS query_expression
주어진 쿼리 표현식으로 뷰를 생성해요. 카탈로그에 같은 이름의 뷰가 이미 존재하면 예외가 발생해요.
TEMPORARY: catalog와 database 네임스페이스를 가진 임시 뷰를 생성하고 뷰를 덮어써요.IF NOT EXISTS: 뷰가 이미 존재하면 아무 일도 일어나지 않아요.
CREATE FUNCTION
CREATE [TEMPORARY|TEMPORARY SYSTEM] FUNCTION
[IF NOT EXISTS] [catalog_name.][db_name.]function_name
AS identifier [LANGUAGE JAVA|SCALA|PYTHON]
[USING [JAR|ARTIFACT] '<path_to_filename>.jar' [, JAR '<path_to_filename>.jar']* ]
[WITH (key1=val1, key2=val2, ...)]
식별자와 선택적 언어 태그로 catalog와 database 네임스페이스를 가진 카탈로그 함수를 생성해요. 카탈로그에 같은 이름의 함수가 이미 존재하면 예외가 발생해요.
언어 태그가 JAVA/SCALA이면 식별자는 UDF의 전체 클래스패스예요. Java/Scala UDF의 구현은 User-defined Functions를 참고하세요.
언어 태그가 PYTHON이면 식별자는 UDF의 정규화된 이름(예: pyflink.table.tests.test_udf.add)이에요. Python UDF의 구현은 Python UDFs를 참고하세요.
언어 태그가 PYTHON인데 현재 프로그램이 Java/Scala 또는 순수 SQL로 작성된 경우에는 Python 의존성을 구성해야 해요.
TEMPORARY: catalog와 database 네임스페이스를 가진 임시 카탈로그 함수를 생성하고 카탈로그 함수를 덮어써요.TEMPORARY SYSTEM: 네임스페이스가 없고 내장 함수를 덮어쓰는 임시 시스템 함수를 생성해요.IF NOT EXISTS: 함수가 이미 존재하면 아무 일도 일어나지 않아요.LANGUAGE JAVA|SCALA|PYTHON: Flink 런타임이 함수를 어떻게 실행할지 지시하는 언어 태그. 현재는 JAVA, SCALA, PYTHON만 지원되며 함수의 기본 언어는 JAVA예요.USING: 함수의 구현과 그 의존성을 포함하는 jar 리소스 목록을 지정해요. jar는 Flink가 현재 지원하는 hdfs/s3/oss 같은 로컬 또는 원격 파일 시스템에 있어야 해요.
주의: 현재 USING 절은 JAVA, SCALA 언어만 지원해요.
CREATE MODEL
CREATE [TEMPORARY] MODEL [IF NOT EXISTS] [catalog_name.][db_name.]model_name
[(
{ <input_column_definition> }[ , ...n]
{ <output_column_definition> }[ , ...n]
)]
[COMMENT model_comment]
WITH (key1=val1, key2=val2, ...)
<input_column_definition>:
column_name column_type [COMMENT column_comment]
<output_column_definition>:
column_name column_type [COMMENT column_comment]
선택적 입력·출력 컬럼 정의로 모델을 생성해요. 카탈로그에 같은 이름의 모델이 이미 존재하면 예외가 발생해요.
TEMPORARY: catalog와 database 네임스페이스를 가진 임시 모델을 생성하고 모델을 덮어써요.IF NOT EXISTS: 모델이 이미 존재하면 아무 일도 일어나지 않아요.Input/Output Columns: 입력 컬럼은 모델 추론에 사용될 특징(features)을 정의해요. 출력 컬럼은 모델이 생성할 예측(predictions)을 정의해요. 각 컬럼은 이름과 데이터 유형을 가져야 해요.WITH OPTIONS: 이 모델과 관련된 추가 정보를 저장하는 데 사용되는 모델 속성이에요. 이 속성들은 보통 기본 모델 제공자를 찾고 생성하는 데 사용돼요.key1=val1표현식의 key와 value 모두 문자열 리터럴이어야 해요.
참고: 모델 속성과 지원되는 모델 유형은 기본 모델 제공자에 따라 다를 수 있어요.
예시 (Examples)
다음 예시는 CREATE MODEL 문의 사용법을 보여줘요.
CREATE MODEL sentiment_analysis_model
INPUT (text STRING COMMENT 'Input text for sentiment analysis')
OUTPUT (sentiment STRING COMMENT 'Predicted sentiment (positive/negative/neutral/mixed)')
COMMENT 'A model for sentiment analysis of text'
WITH (
'provider' = 'openai',
'endpoint' = 'https://api.openai.com/v1/chat/completions',
'api-key' = '<YOUR KEY>',
'model'='gpt-3.5-turbo',
'system-prompt' = 'Classify the text below into one of the following labels: [positive, negative, neutral, mixed]. Output only the label.'
);
CREATE MODEL triton_text_classifier
INPUT (input STRING COMMENT 'Input text for classification')
OUTPUT (output STRING COMMENT 'Classification result')
COMMENT 'A Triton-based text classification model'
WITH (
'provider' = 'triton',
'endpoint' = 'http://localhost:8000/v2/models',
'model-name' = 'text-classification',
'model-version' = '1',
'timeout' = '10000',
'max-retries' = '3'
);