쿼리

쿼리 (Queries)

SELECT 문과 VALUES 문은 TableEnvironmentsqlQuery() 메서드로 지정돼요. 이 메서드는 SELECT 문(또는 VALUES 문)의 결과를 Table로 반환해요. Table후속 SQL 및 Table API 쿼리에서 사용되거나, DataStream으로 변환되거나, TableSink로 작성될 수 있어요. SQL과 Table API 쿼리는 매끄럽게 혼합될 수 있으며 전체적으로 최적화되어 단일 프로그램으로 변환돼요.

출처: 문서

본문

SQL 쿼리에서 테이블에 접근하려면 TableEnvironment에 등록되어야 해요. 테이블은 TableSource, Table, CREATE TABLE 문, DataStream에서 등록될 수 있어요. 또는 사용자는 데이터 소스의 위치를 지정하기 위해 TableEnvironment에 카탈로그를 등록할 수도 있어요.

편의를 위해 Table.toString()은 자동으로 테이블을 해당 TableEnvironment의 고유 이름으로 등록하고 그 이름을 반환해요. 따라서 Table 객체는 아래 예시처럼 SQL 쿼리에 직접 인라인될 수 있어요.

참고: 지원되지 않는 SQL 기능을 포함하는 쿼리는 TableException을 발생시켜요. 배치 및 스트리밍 테이블에서 SQL이 지원하는 기능은 다음 섹션들에 나열돼요.

쿼리 지정하기 (Specifying a Query)

다음 예시들은 등록된 테이블과 인라인된 테이블에 SQL 쿼리를 지정하는 방법을 보여줘요.

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);

// ingest a DataStream from an external source
DataStream<Tuple3<Long, String, Integer>> ds = env.addSource(...);

// SQL query with an inlined (unregistered) table
Table table = tableEnv.fromDataStream(ds, $("user"), $("product"), $("amount"));
Table result = tableEnv.sqlQuery(
  "SELECT SUM(amount) FROM " + table + " WHERE product LIKE '%Rubber%'");

// SQL query with a registered table
// register the DataStream as view "Orders"
tableEnv.createTemporaryView("Orders", ds, $("user"), $("product"), $("amount"));
// run a SQL query on the Table and retrieve the result as a new Table
Table result2 = tableEnv.sqlQuery(
  "SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'");

// create and register a TableSink
final Schema schema = Schema.newBuilder()
    .column("product", DataTypes.STRING())
    .column("amount", DataTypes.INT())
    .build();

final TableDescriptor sinkDescriptor = TableDescriptor.forConnector("filesystem")
    .schema(schema)
    .option("path", "/path/to/file")    
    .format(FormatDescriptor.forFormat("csv")
        .option("field-delimiter", ",")
        .build())
    .build();

tableEnv.createTemporaryTable("RubberOrders", sinkDescriptor);

// 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 env = StreamExecutionEnvironment.getExecutionEnvironment
val tableEnv = StreamTableEnvironment.create(env)

// read a DataStream from an external source
val ds: DataStream[(Long, String, Integer)] = env.addSource(...)

// SQL query with an inlined (unregistered) table
val table = ds.toTable(tableEnv, $"user", $"product", $"amount")
val result = tableEnv.sqlQuery(
  s"SELECT SUM(amount) FROM $table WHERE product LIKE '%Rubber%'")

// SQL query with a registered table
// register the DataStream under the name "Orders"
tableEnv.createTemporaryView("Orders", ds, $"user", $"product", $"amount")
// run a SQL query on the Table and retrieve the result as a new Table
val result2 = tableEnv.sqlQuery(
  "SELECT product, amount FROM Orders WHERE product LIKE '%Rubber%'")

// create and register a TableSink
val schema = Schema.newBuilder()
  .column("product", DataTypes.STRING())
  .column("amount", DataTypes.INT())
  .build()

val sinkDescriptor = TableDescriptor.forConnector("filesystem")
  .schema(schema)
  .format(FormatDescriptor.forFormat("csv")
    .option("field-delimiter", ",")
    .build())
  .build()

tableEnv.createTemporaryTable("RubberOrders", sinkDescriptor)

// 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%'")
env = StreamExecutionEnvironment.get_execution_environment()
table_env = StreamTableEnvironment.create(env)

# SQL query with an inlined (unregistered) table
# elements data type: BIGINT, STRING, BIGINT
table = table_env.from_elements(..., ['user', 'product', 'amount'])
result = table_env \
    .sql_query("SELECT SUM(amount) FROM %s WHERE product LIKE '%%Rubber%%'" % table)

# create and register a TableSink
schema = Schema.new_builder()
    .column("product", DataTypes.STRING())
    .column("amount", DataTypes.INT())
    .build()

sink_descriptor = TableDescriptor.for_connector("filesystem")
    .schema(schema)
    .format(FormatDescriptor.for_format("csv")
        .option("field-delimiter", ",")
        .build())
    .build()

t_env.create_temporary_table("RubberOrders", sink_descriptor)

# 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%'")

쿼리 실행하기 (Execute a Query)

SELECT 문 또는 VALUES 문은 TableEnvironment.executeSql() 메서드를 통해 내용을 로컬로 수집하도록 실행할 수 있어요. 이 메서드는 SELECT 문(또는 VALUES 문)의 결과를 TableResult로 반환해요. SELECT 문과 유사하게 Table 객체는 Table.execute() 메서드를 사용해 쿼리 내용을 로컬 클라이언트로 수집하도록 실행할 수 있어요. TableResult.collect() 메서드는 닫을 수 있는(closeable) 행 반복자(iterator)를 반환해요. 모든 결과 데이터가 수집될 때까지 select 작업은 끝나지 않아요. 리소스 누수를 피하기 위해 CloseableIterator#close() 메서드로 작업을 적극적으로 닫아야 해요. 또한 TableResult.print() 메서드로 select 결과를 클라이언트 콘솔에 출력할 수도 있어요. TableResult의 결과 데이터는 한 번만 접근할 수 있어요. 따라서 collect()print()를 연속해서 호출하면 안 돼요.

TableResult.collect()TableResult.print()는 서로 다른 체크포인팅 설정에서 약간 다른 동작을 해요(스트리밍 작업의 체크포인팅 활성화는 체크포인팅 구성 참고).

  • 체크포인팅이 없는 배치 작업이나 스트리밍 작업의 경우, TableResult.collect()TableResult.print()는 정확히 한 번이나 최소 한 번 보장을 모두 가지지 않아요. 쿼리 결과는 생성되면 클라이언트가 즉시 접근할 수 있지만, 작업이 실패하고 재시작하면 예외가 던져져요.
  • 정확히 한 번 체크포인팅을 가진 스트리밍 작업의 경우, TableResult.collect()TableResult.print()는 종단 간 정확히 한 번 레코드 전달을 보장해요. 결과는 해당 체크포인트가 완료된 후에만 클라이언트가 접근할 수 있어요.
  • 최소 한 번 체크포인팅을 가진 스트리밍 작업의 경우, TableResult.collect()TableResult.print()는 종단 간 최소 한 번 레코드 전달을 보장해요. 쿼리 결과는 생성되면 클라이언트가 즉시 접근할 수 있지만, 동일한 결과가 여러 번 전달될 수 있어요.
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env, settings);

tableEnv.executeSql("CREATE TABLE Orders (`user` BIGINT, product STRING, amount INT) WITH (...)");

// execute SELECT statement
TableResult tableResult1 = tableEnv.executeSql("SELECT * FROM Orders");
// use try-with-resources statement to make sure the iterator will be closed automatically
try (CloseableIterator<Row> it = tableResult1.collect()) {
    while(it.hasNext()) {
        Row row = it.next();
        // handle row
    }
}

// execute Table
TableResult tableResult2 = tableEnv.sqlQuery("SELECT * FROM Orders").execute();
tableResult2.print();

문법 (Syntax)

Flink는 표준 ANSI SQL을 지원하는 Apache Calcite를 사용해 SQL을 파싱해요.

다음 BNF 문법은 배치 및 스트리밍 쿼리에서 지원되는 SQL 기능의 상위 집합을 설명해요. Operations 섹션은 지원되는 기능에 대한 예시를 보여주고 어떤 기능이 배치 또는 스트리밍 쿼리에서만 지원되는지를 나타내요.

query:
    values
  | WITH withItem [ , withItem ]* query
  | {
        select
      | selectWithoutFrom
      | query UNION [ ALL ] query
      | query EXCEPT query
      | query INTERSECT query
    }
    [ ORDER BY orderItem [, orderItem ]* ]
    [ LIMIT { count | ALL } ]
    [ OFFSET start { ROW | ROWS } ]
    [ FETCH { FIRST | NEXT } [ count ] { ROW | ROWS } ONLY]

withItem:
    name
    [ '(' column [, column ]* ')' ]
    AS '(' query ')'

orderItem:
    expression [ ASC | DESC ]

select:
    SELECT [ ALL | DISTINCT ]
    { * | projectItem [, projectItem ]* }
    FROM tableExpression
    [ WHERE booleanExpression ]
    [ GROUP BY { groupItem [, groupItem ]* } ]
    [ HAVING booleanExpression ]
    [ WINDOW windowName AS windowSpec [, windowName AS windowSpec ]* ]

selectWithoutFrom:
    SELECT [ ALL | DISTINCT ]
    { * | projectItem [, projectItem ]* }

projectItem:
    expression [ [ AS ] columnAlias ]
  | tableAlias . *

tableExpression:
    tableReference [, tableReference ]*
  | tableExpression [ NATURAL ] [ LEFT | RIGHT | FULL ] JOIN tableExpression [ joinCondition ]

joinCondition:
    ON booleanExpression
  | USING '(' column [, column ]* ')'

tableReference:
    tablePrimary
    [ matchRecognize ]
    [ [ AS ] alias [ '(' columnAlias [, columnAlias ]* ')' ] ]

tablePrimary:
    [ TABLE ] tablePath [ dynamicTableOptions ] [systemTimePeriod] [[AS] correlationName]
  | LATERAL TABLE '(' functionName '(' expression [, expression ]* ')' ')'
  | [ LATERAL ] '(' query ')'
  | UNNEST '(' expression ')'

tablePath:
    [ [ catalogName . ] databaseName . ] tableName

systemTimePeriod:
    FOR SYSTEM_TIME AS OF dateTimeExpression

dynamicTableOptions:
    /*+ OPTIONS(key=val [, key=val]*) */

key:
    stringLiteral

val:
    stringLiteral

values:
    VALUES expression [, expression ]*

groupItem:
    expression
  | '(' ')'
  | '(' expression [, expression ]* ')'
  | CUBE '(' expression [, expression ]* ')'
  | ROLLUP '(' expression [, expression ]* ')'
  | GROUPING SETS '(' groupItem [, groupItem ]* ')'

windowRef:
    windowName
  | windowSpec

windowSpec:
    [ windowName ]
    '('
    [ ORDER BY orderItem [, orderItem ]* ]
    [ PARTITION BY expression [, expression ]* ]
    [
        RANGE numericOrIntervalExpression {PRECEDING}
      | ROWS numericExpression {PRECEDING}
    ]
    ')'

matchRecognize:
    MATCH_RECOGNIZE '('
    [ PARTITION BY expression [, expression ]* ]
    [ ORDER BY orderItem [, orderItem ]* ]
    [ MEASURES measureColumn [, measureColumn ]* ]
    [ ONE ROW PER MATCH ]
    [ AFTER MATCH
      ( SKIP TO NEXT ROW
      | SKIP PAST LAST ROW
      | SKIP TO FIRST variable
      | SKIP TO LAST variable
      | SKIP TO variable )
    ]
    PATTERN '(' pattern ')'
    [ WITHIN intervalLiteral ]
    DEFINE variable AS condition [, variable AS condition ]*
    ')'

measureColumn:
    expression AS alias

pattern:
    patternTerm [ '|' patternTerm ]*

patternTerm:
    patternFactor [ patternFactor ]*

patternFactor:
    variable [ patternQuantifier ]

patternQuantifier:
    '*'
  | '*?'
  | '+'
  | '+?'
  | '?'
  | '??'
  | '{' { [ minRepeat ], [ maxRepeat ] } '}' ['?']
  | '{' repeat '}'

Flink SQL은 Java와 유사한 식별자(테이블, 속성, 함수 이름)에 대한 어휘 정책을 사용해요:

  • 인용 여부와 관계없이 식별자의 대소문자가 보존돼요.
  • 그 후 식별자는 대소문자 구분(case-sensitively)하여 일치해요.
  • Java와 달리 백틱은 식별자에 영숫자가 아닌 문자를 포함할 수 있게 해줘요 (예: SELECT a AS `my field` FROM t).

문자열 리터럴은 작은따옴표로 감싸야 해요 (예: SELECT 'Hello World'). 이스케이프를 위해 작은따옴표를 두 번 사용해요 (예: SELECT 'It''s me').

Flink SQL> SELECT 'Hello World', 'It''s me';
+-------------+---------+
|      EXPR$0 |  EXPR$1 |
+-------------+---------+
| Hello World | It's me |
+-------------+---------+
1 row in set

문자열 리터럴에서 유니코드 문자를 지원해요. 명시적 유니코드 코드 포인트가 필요하면 다음 문법을 사용해요:

  • 이스케이프 문자로 백슬래시(\)를 사용 (기본값): SELECT U&'\263A'
  • 커스텀 이스케이프 문자 사용: SELECT U&'#263A' UESCAPE '#'

Flink 2.0부터 C 스타일 이스케이프를 사용할 수 있어요.

백슬래시 이스케이프 시퀀스 해석
\b backspace
\f form feed
\n newline
\r carriage return

더 알아보기 (Learn more)