UPDATE 문
UPDATE 문 (UPDATE Statements)
UPDATE 문은 제공된 필터가 있으면 그에 따라 대상 테이블의 행 수준 업데이트를 수행하는 데 사용됩니다.
출처: 문서
본문
주의: 현재
UPDATE문은 배치 모드에서만 지원되며, 대상 테이블 커넥터가 행 수준 업데이트를 지원하기 위해 SupportsRowLevelUpdate 인터페이스를 구현해야 합니다. 관련 인터페이스를 구현하지 않은 테이블에UPDATE를 시도하면 예외가 발생합니다. 현재 Flink가 유지관리하는 커넥터 중 UPDATE를 지원하는 커넥터는 아직 없습니다.
UPDATE 문 실행 (Run an UPDATE statement)
Java:
UPDATE 문은 TableEnvironment의 executeSql() 메서드로 실행할 수 있습니다. executeSql()은 Flink 작업을 즉시 제출하고, 제출된 작업과 연결된 TableResult 인스턴스를 반환합니다. 다음 예제는 TableEnvironment에서 단일 UPDATE 문을 실행하는 방법을 보여줍니다.
EnvironmentSettings settings = EnvironmentSettings.newInstance().inBatchMode().build();
TableEnvironment tEnv = TableEnvironment.create(settings);
// "Orders"라는 테이블 등록
tEnv.executeSql("CREATE TABLE Orders (`user` STRING, product STRING, amount INT) WITH (...)");
// 값 삽입
tEnv.executeSql("insert into Orders values ('Lili', 'Apple', 1), ('Jessica', 'Banana', 1)").await();
tEnv.executeSql("SELECT * FROM Orders").print();
// amount 전체 업데이트
tEnv.executeSql("UPDATE Orders SET `amount` = `amount` * 2").await();
tEnv.executeSql("SELECT * FROM Orders").print();
// 필터로 업데이트
tEnv.executeSql("UPDATE Orders SET `product` = 'Orange' WHERE `user` = 'Lili'").await();
tEnv.executeSql("SELECT * FROM Orders").print();
Scala:
val env = StreamExecutionEnvironment.getExecutionEnvironment()
val settings = EnvironmentSettings.newInstance().inBatchMode().build()
val tEnv = StreamTableEnvironment.create(env, settings)
// "Orders"라는 테이블 등록
tEnv.executeSql("CREATE TABLE Orders (`user` STRING, product STRING, amount INT) WITH (...)");
// 값 삽입
tEnv.executeSql("insert into Orders values ('Lili', 'Apple', 1), ('Jessica', 'Banana', 1)").await();
tEnv.executeSql("SELECT * FROM Orders").print();
// amount 전체 업데이트
tEnv.executeSql("UPDATE Orders SET `amount` = `amount` * 2").await();
tEnv.executeSql("SELECT * FROM Orders").print();
// 필터로 업데이트
tEnv.executeSql("UPDATE Orders SET `product` = 'Orange' WHERE `user` = 'Lili'").await();
tEnv.executeSql("SELECT * FROM Orders").print();
Python:
env_settings = EnvironmentSettings.in_batch_mode()
table_env = TableEnvironment.create(env_settings)
# "Orders"라는 테이블 등록
table_env.executeSql("CREATE TABLE Orders (`user` STRING, product STRING, amount INT) WITH (...)");
# 값 삽입
table_env.executeSql("insert into Orders values ('Lili', 'Apple', 1), ('Jessica', 'Banana', 1)").wait();
table_env.executeSql("SELECT * FROM Orders").print();
# amount 전체 업데이트
table_env.executeSql("UPDATE Orders SET `amount` = `amount` * 2").wait();
table_env.executeSql("SELECT * FROM Orders").print();
# 필터로 업데이트
table_env.executeSql("UPDATE Orders SET `product` = 'Orange' WHERE `user` = 'Lili'").wait();
table_env.executeSql("SELECT * FROM Orders").print();
SQL CLI:
UPDATE 문은 SQL CLI에서 실행할 수 있습니다.
Flink SQL> SET 'execution.runtime-mode' = 'batch';
[INFO] Session property has been set.
Flink SQL> CREATE TABLE Orders (`user` STRING, product STRING, amount INT) with (...);
[INFO] Execute statement succeeded.
Flink SQL> INSERT INTO Orders VALUES ('Lili', 'Apple', 1), ('Jessica', 'Banana', 1);
[INFO] Submitting SQL update statement to the cluster...
[INFO] SQL update statement has been successfully submitted to the cluster:
Job ID: bd2c46a7b2769d5c559abd73ecde82e9
Flink SQL> SELECT * FROM Orders;
Flink SQL> UPDATE Orders SET amount = 2;
UPDATE ROWS
UPDATE [catalog_name.][db_name.]table_name SET column_name1 = expression1 [, column_name2 = expression2, ...][ WHERE condition ]