INSERT Statement
INSERT Statement (INSERT 문)
INSERT 문은 테이블에 행을 추가하는 데 사용돼요. INSERT INTO나 INSERT OVERWRITE를 사용해 select 쿼리 또는 값 리스트에서 행을 삽입할 수 있어요.
출처: INSERT Statement
본문
INSERT 문은 테이블에 행을 추가하는 데 사용돼요.
Run an INSERT statement (INSERT 문 실행)
단일 INSERT 문은 TableEnvironment의 executeSql() 메서드를 통해 실행할 수 있어요. INSERT 문을 위한 executeSql() 메서드는 즉시 Flink 잡을 제출하고, 제출된 잡과 연결된 TableResult 인스턴스를 반환해요.
여러 INSERT 문은 TableEnvironment.createStatementSet() 메서드로 만들 수 있는 StatementSet의 addInsertSql() 메서드를 통해 실행할 수 있어요. addInsertSql() 메서드는 지연 실행이며, StatementSet.execute()가 호출될 때만 실행돼요.
다음 예제는 TableEnvironment에서 단일 INSERT 문을 실행하고, StatementSet에서 여러 INSERT 문을 실행하는 방법을 보여줘요.
TableEnvironment tEnv = TableEnvironment.create(...);
// register a source table named "Orders" and a sink table named "RubberOrders"
tEnv.executeSql("CREATE TABLE Orders (`user` BIGINT, product VARCHAR, amount INT) WITH (...)");
tEnv.executeSql("CREATE TABLE RubberOrders(product VARCHAR, amount INT) WITH (...)");
// run a single INSERT query on the registered source table and emit the result to registered sink table
TableResult tableResult1 = tEnv.executeSql(
"INSERT INTO RubberOrders SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'");
// get job status through TableResult
System.out.println(tableResult1.getJobClient().get().getJobStatus());
//----------------------------------------------------------------------------
// register another sink table named "GlassOrders" for multiple INSERT queries
tEnv.executeSql("CREATE TABLE GlassOrders(product VARCHAR, amount INT) WITH (...)");
// run multiple INSERT queries on the registered source table and emit the result to registered sink tables
StatementSet stmtSet = tEnv.createStatementSet();
// only single INSERT query can be accepted by `addInsertSql` method
stmtSet.addInsertSql(
"INSERT INTO RubberOrders SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'");
stmtSet.addInsertSql(
"INSERT INTO GlassOrders SELECT product, amount FROM Orders WHERE product LIKE '%Glass%'");
// execute all statements together
TableResult tableResult2 = stmtSet.execute();
// get job status through TableResult
System.out.println(tableResult2.getJobClient().get().getJobStatus());
val tEnv = TableEnvironment.create(...)
// register a source table named "Orders" and a sink table named "RubberOrders"
tEnv.executeSql("CREATE TABLE Orders (`user` BIGINT, product STRING, amount INT) WITH (...)")
tEnv.executeSql("CREATE TABLE RubberOrders(product STRING, amount INT) WITH (...)")
// run a single INSERT query on the registered source table and emit the result to registered sink table
val tableResult1 = tEnv.executeSql(
"INSERT INTO RubberOrders SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'")
// get job status through TableResult
println(tableResult1.getJobClient().get().getJobStatus())
//----------------------------------------------------------------------------
// register another sink table named "GlassOrders" for multiple INSERT queries
tEnv.executeSql("CREATE TABLE GlassOrders(product VARCHAR, amount INT) WITH (...)")
// run multiple INSERT queries on the registered source table and emit the result to registered sink tables
val stmtSet = tEnv.createStatementSet()
// only single INSERT query can be accepted by `addInsertSql` method
stmtSet.addInsertSql(
"INSERT INTO RubberOrders SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'")
stmtSet.addInsertSql(
"INSERT INTO GlassOrders SELECT product, amount FROM Orders WHERE product LIKE '%Glass%'")
// execute all statements together
val tableResult2 = stmtSet.execute()
// get job status through TableResult
println(tableResult2.getJobClient().get().getJobStatus())
table_env = TableEnvironment.create(...)
# register a source table named "Orders" and a sink table named "RubberOrders"
table_env.execute_sql("CREATE TABLE Orders (`user` BIGINT, product STRING, amount INT) WITH (...)")
table_env.execute_sql("CREATE TABLE RubberOrders(product STRING, amount INT) WITH (...)")
# run a single INSERT query on the registered source table and emit the result to registered sink table
table_result1 = table_env \
.execute_sql("INSERT INTO RubberOrders SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'")
# get job status through TableResult
print(table_result1get_job_client().get_job_status())
#----------------------------------------------------------------------------
# register another sink table named "GlassOrders" for multiple INSERT queries
table_env.execute_sql("CREATE TABLE GlassOrders(product VARCHAR, amount INT) WITH (...)")
# run multiple INSERT queries on the registered source table and emit the result to registered sink tables
stmt_set = table_env.create_statement_set()
# only single INSERT query can be accepted by `add_insert_sql` method
stmt_set \
.add_insert_sql("INSERT INTO RubberOrders SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'")
stmt_set \
.add_insert_sql("INSERT INTO GlassOrders SELECT product, amount FROM Orders WHERE product LIKE '%Glass%'")
# execute all statements together
table_result2 = stmt_set.execute()
# get job status through TableResult
print(table_result2.get_job_client().get_job_status())
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 (...);
Flink SQL> SHOW TABLES;
Orders
RubberOrders
Flink SQL> INSERT INTO RubberOrders SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%';
[INFO] Submitting SQL update statement to the cluster...
[INFO] Table update statement has been successfully submitted to the cluster:
Insert from select queries (select 쿼리에서 삽입)
쿼리 결과는 insert 절을 사용하여 테이블에 삽입할 수 있어요.
Syntax
[EXECUTE] INSERT { INTO | OVERWRITE } [catalog_name.][db_name.]table_name [PARTITION part_spec] [column_list] select_statement
part_spec:
(part_col_name1=val1 [, part_col_name2=val2, ...])
column_list:
(col_name1 [, column_name2, ...])
OVERWRITE
INSERT OVERWRITE는 테이블 또는 파티션의 기존 데이터를 모두 덮어써요. 그렇지 않으면 새 데이터가 추가돼요.
PARTITION
PARTITION 절은 이 삽입의 정적 파티션 열을 포함해야 해요.
COLUMN LIST
테이블 T(a INT, b INT, c INT)가 주어지면 Flink는 INSERT INTO T(c, b) SELECT x, y FROM S를 지원해요. 기대는 'x'가 열 'c'에, 'y'가 열 'b'에 쓰이고 'a'는 NULL로 설정된다는 것이에요 (열 'a'가 nullable이라고 가정).
부분 열 업데이트를 처리할 때 비대상 열을 null 값으로 덮어쓰는 것을 피하고 싶은 커넥터 개발자는 사용자의 insert 문이 지정한 대상 열 정보를 {{< gh_link file="flink-table/flink-table-common/src/main/java/org/apache/flink/table/connector/sink/DynamicTableSink.java" name="DynamicTableSink$Context.getTargetColumns()" >}}에서 얻어 부분 업데이트를 처리할 방법을 결정할 수 있어요.
Examples
-- Creates a partitioned table
CREATE TABLE country_page_view (user STRING, cnt INT, date STRING, country STRING)
PARTITIONED BY (date, country)
WITH (...)
-- Appends rows into the static partition (date='2019-8-30', country='China')
INSERT INTO country_page_view PARTITION (date='2019-8-30', country='China')
SELECT user, cnt FROM page_view_source;
-- Key word EXECUTE can be added at the beginning of Insert to indicate explicitly that we are going to execute the statement,
-- it is equivalent to Statement without the key word.
EXECUTE INSERT INTO country_page_view PARTITION (date='2019-8-30', country='China')
SELECT user, cnt FROM page_view_source;
-- Appends rows into partition (date, country), where date is static partition with value '2019-8-30',
-- country is dynamic partition whose value is dynamic determined by each row.
INSERT INTO country_page_view PARTITION (date='2019-8-30')
SELECT user, cnt, country FROM page_view_source;
-- Overwrites rows into static partition (date='2019-8-30', country='China')
INSERT OVERWRITE country_page_view PARTITION (date='2019-8-30', country='China')
SELECT user, cnt FROM page_view_source;
-- Overwrites rows into partition (date, country), where date is static partition with value '2019-8-30',
-- country is dynamic partition whose value is dynamic determined by each row.
INSERT OVERWRITE country_page_view PARTITION (date='2019-8-30')
SELECT user, cnt, country FROM page_view_source;
-- Appends rows into the static partition (date='2019-8-30', country='China')
-- the column cnt is set to NULL
INSERT INTO country_page_view PARTITION (date='2019-8-30', country='China') (user)
SELECT user FROM page_view_source;
Insert values into tables (테이블에 값 삽입)
INSERT...VALUES 문은 SQL에서 직접 테이블에 데이터를 삽입하는 데 사용할 수 있어요.
Syntax
[EXECUTE] INSERT { INTO | OVERWRITE } [catalog_name.][db_name.]table_name VALUES values_row [, values_row ...]
values_row:
(val1 [, val2, ...])
OVERWRITE
INSERT OVERWRITE는 테이블의 기존 데이터를 모두 덮어써요. 그렇지 않으면 새 데이터가 추가돼요.
Examples
CREATE TABLE students (name STRING, age INT, gpa DECIMAL(3, 2)) WITH (...);
INSERT INTO students
VALUES ('fred flintstone', 35, 1.28), ('barney rubble', 32, 2.32);
Insert into multiple tables (여러 테이블에 삽입)
STATEMENT SET을 사용하면 하나의 문에서 여러 테이블에 데이터를 삽입할 수 있어요.
Syntax
EXECUTE STATEMENT SET
BEGIN
insert_statement;
...
insert_statement;
END;
insert_statement:
<insert_from_select>|<insert_from_values>
Examples
CREATE TABLE students (name STRING, age INT, gpa DECIMAL(3, 2)) WITH (...);
EXECUTE STATEMENT SET
BEGIN
INSERT INTO students
VALUES ('fred flintstone', 35, 1.28), ('barney rubble', 32, 2.32);
INSERT INTO students
VALUES ('fred flintstone', 35, 1.28), ('barney rubble', 32, 2.32);
END;
ON CONFLICT clause (ON CONFLICT 절)
쿼리가 싱크 테이블의 기본 키와 다른 upsert 키를 가진 업데이트 테이블을 만들 때, 서로 다른 upsert 키를 가진 여러 레코드가 같은 기본 키에 매핑될 수 있어요. ON CONFLICT 절은 싱크에서 이러한 기본 키 충돌을 해결하는 방법을 지정해요.
ON CONFLICT가 필요한 경우
기본적으로 쿼리의 upsert 키가 싱크 테이블의 기본 키와 다를 때마다 Flink는 명시적 ON CONFLICT 절을 요구해요. 그것이 없으면 쿼리는 플래닝 시간에 실패해요. 이는 쿼리에 실제 충돌 시나리오가 있는지, 아니면 논리 문제(예: 누락된 GROUP BY)가 있는지 고려하도록 강제해요.
이 검사는 구성 옵션 table.exec.sink.require-on-conflict(기본값: true)에 의해 제어돼요. false로 설정하면 ON CONFLICT 절이 필요 없는 레거시 동작으로 복원되지만, 비결정적 결과를 초래할 수 있어요.
또는 충돌 키에 대한 일관성 보장이 필요 없다면 table.exec.sink.upsert-materialize를 NONE으로 설정하여 싱크 upsert materializer를 완전히 비활성화할 수 있어요. 이렇게 하면 materializer 연산자가 파이프라인에서 제거되어 버퍼링, 압축, 충돌 해결이 수행되지 않아요. 레코드는 도착하는 대로 싱크에 직접 전달돼요.
Syntax
[EXECUTE] INSERT INTO [catalog_name.][db_name.]table_name
select_statement
ON CONFLICT conflict_action
conflict_action:
DO NOTHING
| DO ERROR
| DO DEDUPLICATE
Strategies (전략)
DO ERROR
서로 다른 upsert 키를 가진 여러 레코드가 같은 기본 키에 매핑되면 런타임에 예외를 발생시켜요. 실제 충돌이 없다고 믿을 때 사용해요 — 예를 들어 플래너가 upsert 키가 기본 키와 일치함을 증명할 수는 없지만 논리적으로 동등하다는 것을 알 때요.
버퍼링된 레코드는 충돌 검사 전에 워터마크 진행 시 압축되므로, changelog 재정렬로 인한 일시적 무질서가 거짓 오류를 일으키지 않아요.
INSERT INTO product_orders
SELECT p.name, o.order_id
FROM orders o JOIN products p ON o.product_name = p.name
ON CONFLICT DO ERROR;
DO NOTHING
주어진 기본 키에 대해 먼저 도착하는 레코드를 유지하고, 이후에 도착하는 충돌 레코드를 조용히 버려요. 서로 다른 upsert 키에서 중복 기본 키 값을 버리는 것이 허용될 때 사용해요.
DO ERROR처럼 이 전략은 충돌 해결 적용 전에 워터마크 기반 압축을 사용해요.
INSERT INTO product_orders
SELECT p.name, o.order_id
FROM orders o JOIN products p ON o.product_name = p.name
ON CONFLICT DO NOTHING;
DO DEDUPLICATE
{{< hint warning >}}
DO DEDUPLICATE는 철회(retraction) 시 롤백을 지원하기 위해 상태에 기본 키별 변경 사항의 전체 기록을 유지해요. 이로 인해 DO ERROR와 DO NOTHING에 비해 상태 사용량이 상당히 높아져요.
{{< /hint >}}
기본 키별 변경 사항의 전체 기록을 유지하여 철회가 올바르게 롤백될 수 있게 해요. 같은 기본 키에 대한 실제 다중 소스 업데이트가 발생하고 정확성을 희생할 수 없을 때 가장 올바른 전략이에요.
INSERT INTO product_orders
SELECT p.name, o.order_id
FROM orders o JOIN products p ON o.product_name = p.name
ON CONFLICT DO DEDUPLICATE;
충돌이 발생하는 방식
충돌은 쿼리의 upsert 키가 싱크 테이블의 기본 키와 다를 때 발생해요. 예를 들어 조인 조건에서 파생된 upsert 키를 가진 조인 결과가 있지만 대상 테이블에 다른 기본 키가 있다고 생각해보세요. 그러면 서로 다른 상위 upsert 키의 레코드가 싱크에서 같은 기본 키에 충돌할 수 있어요.
철회(-U)와 업데이트(+U) 메시지는 파이프라인을 통해 다른 경로로 이동할 수 있으므로 싱크에 순서가 어긋나 도착할 수 있어요. DO ERROR와 DO NOTHING은 워터마크 기반 압축을 사용해 일관된 변경 집합을 기다린 후 충돌 해결을 적용하므로, 일시적 재정렬로 인한 거짓 긍정을 방지해요.
워터마크 기반 압축
조인과 같은 연산자가 생성한 changelog 메시지는 싱크에 순서가 어긋나 도착할 수 있어요. 한 행에 대한 철회(-U)가 같은 기본 키를 공유하는 다른 행의 새 삽입(+I)보다 나중에 도착하여, 그 키에 대해 두 개의 활성 레코드가 존재하는 것처럼 보일 수 있는데, 이것이 거짓 충돌이에요.
워터마크 기반 압축은 기본 키와 upsert 키로 키가 지정된 들어오는 레코드를 버퍼링하여 이 문제를 해결해요. 워터마크가 진행되면 해당 워터마크까지의 타임스탬프가 있는 모든 버퍼링된 레코드가 압축돼요: 같은 upsert 키에 대한 일치하는 삽입과 철회 쌍은 서로 상쇄돼요 (예: +I와 -D, 또는 -U와 +U 쌍).
예제. 위의 orders JOIN products 쿼리를 사용해, 주문 1이 제품을 Laptop에서 Phone으로 바꾸고 주문 3도 Laptop용이라고 가정해보세요. 조인은 다음 changelog 레코드를 방출해요:
+I[Laptop, 1] -- upsert key: order_id=1
+I[Laptop, 3] -- upsert key: order_id=3
-U[Laptop, 1] -- upsert key: order_id=1 (retraction for order 1's old product)
+U[Phone, 1] -- upsert key: order_id=1 (order 1 now maps to Phone)
압축 없이 처음 두 +I 레코드가 도착한 후 연산자는 PK Laptop에 대해 서로 다른 upsert 키(order_id=1과 order_id=3)를 가진 두 개의 활성 레코드를 보게 되는데, 이것이 거짓 충돌이에요. 압축을 사용하면 연산자가 워터마크를 기다려요. 그런 다음 철회 -U[Laptop, 1]이 이전의 +I[Laptop, 1](같은 upsert 키 order_id=1)을 상쇄하여 PK Laptop에 +I[Laptop, 3]만 남게 해요 — 충돌이 없어요.
압축 후 기본 키당 0개 또는 1개의 레코드만 남으면 충돌이 없어요. 서로 다른 upsert 키를 가진 여러 레코드가 여전히 남아 있으면 진짜 충돌이 존재하며 선택한 전략(DO ERROR 또는 DO NOTHING)에 의해 해결돼요. DO DEDUPLICATE는 워터마크 기반 압축을 사용하지 않아요; 대신 철회 시 올바른 롤백을 지원하기 위해 상태에 변경 사항의 전체 기록을 유지해요.
Examples
-- Source and dimension tables
CREATE TABLE orders (
order_id BIGINT,
product_name STRING,
quantity INT,
PRIMARY KEY(order_id) NOT ENFORCED
) WITH (...);
CREATE TABLE products (
name STRING,
PRIMARY KEY(name) NOT ENFORCED
) WITH (...);
-- Sink table
CREATE TABLE product_orders (
product_name STRING,
last_order_id BIGINT,
PRIMARY KEY(product_name) NOT ENFORCED
) WITH (...);
-- This join produces an upsert key that may differ from the sink's PK,
-- so ON CONFLICT is required.
INSERT INTO product_orders
SELECT p.name, o.order_id
FROM orders o JOIN products p ON o.product_name = p.name
ON CONFLICT DO NOTHING;
소스 테이블에 다음 데이터가 주어졌다고 가정해보세요:
orders: products:
+----------+--------------+----------+ +--------+
| order_id | product_name | quantity | | name |
+----------+--------------+----------+ +--------+
| 1 | Laptop | 2 | | Laptop |
| 2 | Phone | 1 | | Phone |
| 3 | Laptop | 5 | +--------+
+----------+--------------+----------+
조인은 product_orders에 대해 다음 changelog 레코드를 생성해요:
+I[Laptop, 1] -- upsert key: order_id=1
+I[Phone, 2] -- upsert key: order_id=2
+I[Laptop, 3] -- upsert key: order_id=3 ← conflicts with order_id=1 on PK 'Laptop'
서로 다른 upsert 키(order_id=1과 order_id=3)를 가진 두 레코드가 같은 기본 키(product_name='Laptop')를 대상으로 해요. 이것이 각 전략이 다르게 해결하는 충돌이에요:
-
DO ERROR— 두 개의 서로 다른 upsert 키가 같은 기본 키에 매핑되므로 런타임 예외를 발생시켜요. -
DO NOTHING— 첫 번째 레코드를 유지하고 충돌을 버려요:product_name last_order_id Laptop 1 Phone 2 -
DO DEDUPLICATE— 둘 다 수용해요; 마지막으로 도착한 값이 보여요:product_name last_order_id Laptop 3 Phone 2
철회 시 어떤 일이 일어날까? 나중에 주문 3이 소스에서 삭제되면 조인은 철회 -D[Laptop, 3]를 방출해요:
DO NOTHING— 철회는 효과가 없어요.(Laptop, 3)은 결코 쓰이지 않았으니까요. Laptop 행은last_order_id=1로 남아요.DO DEDUPLICATE— 이전 값으로 롤백돼요. Laptop은 주문 1로 돌아가{(Laptop, 1), (Phone, 2)}를 생성해요. 상태에 유지된 전체 기록이 이 올바른 롤백을 가능하게 해요.