INSERT 문
INSERT 문 (INSERT Statement)
INSERT 문은 테이블에 행을 추가하는 데 사용됩니다.
출처: 문서
본문
INSERT 문은 테이블에 행을 추가하는 데 사용됩니다.
INSERT 문 실행
Java / Scala: 단일 INSERT 문은 TableEnvironment 의 executeSql() 메서드를 통해 실행할 수 있습니다. INSERT 문의 executeSql() 메서드는 Flink 작업을 즉시 제출하고 제출된 작업과 연결된 TableResult 인스턴스를 반환합니다. 여러 INSERT 문은 TableEnvironment.createStatementSet() 메서드로 만들 수 있는 StatementSet 의 addInsertSql() 메서드를 통해 실행할 수 있습니다. addInsertSql() 메서드는 지연 실행이며 StatementSet.execute() 가 호출될 때만 실행됩니다.
Python: 단일 INSERT 문은 TableEnvironment 의 execute_sql() 메서드를 통해 실행할 수 있습니다. INSERT 문의 execute_sql() 메서드는 Flink 작업을 즉시 제출하고 제출된 작업과 연결된 TableResult 인스턴스를 반환합니다. 여러 INSERT 문은 TableEnvironment.create_statement_set() 메서드로 만들 수 있는 StatementSet 의 add_insert_sql() 메서드를 통해 실행할 수 있습니다. add_insert_sql() 메서드는 지연 실행이며 StatementSet.execute() 가 호출될 때만 실행됩니다.
SQL CLI: 단일 INSERT 문은 SQL CLI 에서 실행할 수 있습니다.
Java:
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());
Scala:
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())
Python:
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())
SQL CLI:
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:
Select 쿼리에서 삽입
쿼리 결과는 insert 절을 사용해 테이블에 삽입할 수 있습니다.
문법
[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 문이 지정한 대상 열에 대한 정보를 DynamicTableSink$Context.getTargetColumns() 에서 얻고 부분 업데이트를 처리하는 방법을 결정할 수 있습니다.
예제
-- 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 문은 SQL 에서 직접 테이블에 데이터를 삽입하는 데 사용될 수 있습니다.
문법
[EXECUTE] INSERT { INTO | OVERWRITE } [catalog_name.][db_name.]table_name VALUES values_row [, values_row ...]
values_row:
(val1 [, val2, ...])
OVERWRITE: INSERT OVERWRITE 는 테이블의 기존 데이터를 덮어씁니다. 그렇지 않으면 새 데이터가 추가됩니다.
예제
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);
여러 테이블에 삽입
STATEMENT SET 은 하나의 문으로 여러 테이블에 데이터를 삽입하는 데 사용될 수 있습니다.
문법
EXECUTE STATEMENT SET
BEGIN
insert_statement;
...
insert_statement;
END;
insert_statement:
<insert_from_select>|<insert_from_values>
예제
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 절
쿼리가 sink 테이블의 기본 키와 다른 upsert 키를 가진 업데이트 테이블을 생성할 때, 서로 다른 upsert 키를 가진 여러 레코드가 같은 기본 키에 매핑될 수 있습니다. ON CONFLICT 절은 sink 에서 이러한 기본 키 충돌을 해결하는 방법을 지정합니다.
ON CONFLICT 가 필요한 경우
기본적으로 Flink 는 쿼리의 upsert 키가 sink 테이블의 기본 키와 다를 때마다 명시적인 ON CONFLICT 절을 요구합니다. 없으면 쿼리는 플래닝(planning) 시점에 실패합니다. 이는 쿼리에 실제 충돌 시나리오가 있는지, 아니면 로직 문제(예: 누락된 GROUP BY)가 있는지 고려하도록 강제합니다.
이 검사는 구성 옵션 table.exec.sink.require-on-conflict (기본값: true) 로 제어됩니다. false 로 설정하면 ON CONFLICT 절이 필요하지 않았던 레거시 동작으로 복원되지만 비결정적 결과를 초래할 수 있습니다.
또는 충돌 키에 대한 일관성 보장이 필요하지 않다면 table.exec.sink.upsert-materialize 를 NONE 으로 설정해 sink upsert 구체화기를 완전히 비활성화할 수 있습니다. 이는 파이프라인에서 구체화 operator 를 제거하므로 버퍼링, 압축(compaction), 충돌 해결이 수행되지 않습니다. 레코드는 도착하는 대로 sink 에 직접 전달됩니다.
문법
[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
DO DEDUPLICATE 는 기본 키당 변경의 전체 이력을 상태에 유지해 재시도(retraction) 시 롤백을 지원합니다. 이는 DO ERROR 와 DO NOTHING 에 비해 상태 사용량이 상당히 높아집니다.
재시도가 올바르게 롤백될 수 있도록 기본 키당 변경의 전체 이력을 유지합니다. 같은 기본 키에 대한 실제 다중 소스 업데이트가 발생하고 정확성을 희생할 수 없는 경우 가장 올바른 전략입니다.
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 키가 sink 테이블의 기본 키와 다를 때 충돌이 발생합니다. 예를 들어 join 결과의 upsert 키가 join 조건에서 파생되었지만 대상 테이블에 다른 기본 키가 있는 경우를 생각해 보세요. 서로 다른 업스트림 upsert 키의 레코드는 sink 에서 같은 기본 키에 충돌할 수 있습니다.
재시도(-U) 와 업데이트(+U) 메시지는 파이프라인을 통해 다른 경로로 이동할 수 있으므로 sink 에 무순서로 도착할 수 있습니다. DO ERROR 와 DO NOTHING 은 워터마크 기반 압축을 사용해 일관된 변경 집합을 기다린 후 충돌 해결을 적용하므로, 일시적인 재정렬로 인한 잘못된 긍정(false positive)을 방지합니다.
워터마크 기반 압축
join 같은 operator 가 생성한 changelog 메시지는 sink 에 무순서로 도착할 수 있습니다. 한 행의 재시도(-U) 가 같은 기본 키를 공유하는 다른 행의 새 삽입(+I) 이후에 도착해 해당 키에 두 개의 활성 레코드가 있는 것처럼 보일 수 있습니다 — 잘못된 충돌입니다.
워터마크 기반 압축은 들어오는 레코드를 기본 키와 upsert 키로 키잉해 버퍼링함으로써 이를 해결합니다. 워터마크가 진행되면 해당 워터마크까지의 타임스탬프를 가진 모든 버퍼된 레코드가 압축됩니다: 같은 upsert 키에 대한 일치하는 삽입과 재시도 쌍은 서로 상쇄됩니다(예: +I 와 -D, 또는 -U 와 +U 쌍).
예제. 위의 orders JOIN products 쿼리를 사용하여, order 1 이 제품을 Laptop 에서 Phone 으로 변경하고 order 3 도 Laptop 용이라고 가정합니다. join 은 다음 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 레코드가 도착한 후 operator 는 서로 다른 upsert 키(order_id=1 과 order_id=3)를 가진 PK Laptop 에 대해 두 개의 활성 레코드를 봅니다 — 잘못된 충돌입니다. 압축을 사용하면 operator 는 워터마크를 기다립니다. 그러면 재시도 -U[Laptop, 1] 이 이전 +I[Laptop, 1] (같은 upsert 키 order_id=1) 을 상쇄해 PK Laptop 에 대해 +I[Laptop, 3] 만 남깁니다 — 충돌이 없습니다.
압축 후 기본 키당 레코드가 0개 또는 1개 남으면 충돌이 없습니다. 서로 다른 upsert 키를 가진 여러 레코드가 여전히 남아 있으면 실제 충돌이 존재하며 선택한 전략(DO ERROR 또는 DO NOTHING) 으로 해결됩니다. DO DEDUPLICATE 는 워터마크 기반 압축을 사용하지 않습니다. 대신 재시도 시 올바른 롤백을 지원하기 위해 상태에 변경의 전체 이력을 유지합니다.
예제
-- 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 | +--------+
+----------+--------------+----------+
join 은 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 |
재시도 시 어떻게 되나요? order 3 이 나중에 소스에서 삭제되면 join 은 재시도 -D[Laptop, 3] 을 내보냅니다:
DO NOTHING—(Laptop, 3)은 결코 기록되지 않았으므로 재시도는 효과가 없습니다. Laptop 행은last_order_id=1로 유지됩니다.DO DEDUPLICATE— 이전 값으로 롤백합니다. Laptop 은 order 1 로 돌아가{(Laptop, 1), (Phone, 2)}를 생성합니다. 상태에 유지된 전체 이력이 이 올바른 롤백을 가능하게 합니다.