스트리밍 개념

스트리밍 개념 (Streaming Concepts)

Flink의 Table API와 SQL 지원은 배치 및 스트림 처리를 위한 통합 API입니다. 즉, Table API와 SQL 쿼리는 입력이 유한 배치 입력이든 무한 스트림 입력이든 동일한 의미론을 가집니다.

다음 페이지에서는 스트리밍 데이터에 대한 Flink 관계형 API의 개념, 실용적 한계, 스트림 특화 구성 파라미터를 설명합니다.

출처: 문서

본문

상태 관리 (State Management)

스트리밍 모드에서 실행되는 테이블 프로그램은 상태 저장 스트림 프로세서로서의 Flink의 모든 능력을 활용합니다.

특히 테이블 프로그램은 상태 크기와 장애 허용에 대한 다양한 요구 사항을 처리하기 위해 상태 백엔드(state backend)와 다양한 체크포인팅 옵션으로 구성할 수 있습니다. 실행 중인 Table API & SQL 파이프라인의 세이브포인트를 찍고 나중에 애플리케이션 상태를 복원하는 것도 가능합니다.

상태 사용 (State Usage)

Table API & SQL 프로그램의 선언적 특성 때문에 파이프라인 내에서 어디에 얼마나 많은 상태가 사용되는지 항상 명확하지 않습니다. 플래너(planner)가 올바른 결과를 계산하는 데 상태가 필요한지 결정합니다. 파이프라인은 주어진 옵티마이저 규칙 집합에서 가능한 한 적은 상태를 차지하도록 최적화됩니다.

개념적으로 소스 테이블은 상태에 완전히 보관되지 않습니다. 구현자는 논리 테이블(즉, 동적 테이블)을 다룹니다. 그 상태 요구 사항은 사용되는 연산에 따라 달라집니다.

상태 저장 연산자 (Stateful Operators)

조인, 집계, 중복 제거 같은 상태 저장 연산을 포함하는 쿼리는 장애 허용 저장소에 중간 결과를 유지해야 하며, 이를 위해 Flink의 상태 추상화가 사용됩니다.

예를 들어 두 테이블의 일반 SQL 조인은 연산자가 두 입력 테이블을 상태에 완전히 유지해야 합니다. 올바른 SQL 의미론을 위해 런타임은 양쪽 어느 시점에서든 매칭이 발생할 수 있다고 가정해야 합니다. Flink는 워터마크 개념을 활용해 상태 크기를 작게 유지하는 것을 목표로 하는 최적화된 윈도우 및 구간 조인(interval join)을 제공합니다.

또 다른 예는 단어 개수를 계산하는 다음 쿼리입니다.

CREATE TABLE doc (
    word STRING
) WITH (
    'connector' = '...'
);
CREATE TABLE word_cnt (
    word STRING PRIMARY KEY NOT ENFORCED,
    cnt  BIGINT
) WITH (
    'connector' = '...'
);

INSERT INTO word_cnt
SELECT word, COUNT(1) AS cnt
FROM doc
GROUP BY word;

word 필드는 그룹핑 키로 사용되며, 연속 쿼리는 관찰하는 각 word에 대한 개수를 싱크에 기록합니다. word 값은 시간에 따라 진화하고, 연속 쿼리는 끝나지 않으므로 프레임워크는 관찰된 각 word 값에 대한 개수를 유지해야 합니다. 결과적으로 쿼리의 총 상태 크기는 더 많은 word 값을 관찰함에 따라 계속 커집니다.

SELECT ... FROM ... WHERE처럼 필드 프로젝션이나 필터만으로 구성된 쿼리는 보통 무상태(stateless) 파이프라인입니다. 그러나 어떤 상황에서는 상태 저장 연산이 입력의 특성(예: 입력이 UPDATE_BEFORE 없이 changelog인 경우, Table to Stream Conversion 참고) 또는 사용자 구성(table-exec-source-cdc-events-duplicate 참고)을 통해 암시적으로 파생됩니다.

다음 그림은 upsert kafka 소스를 조회하는 SELECT ... FROM 문을 보여줍니다.

CREATE TABLE upsert_kafka (
    id INT PRIMARY KEY NOT ENFORCED,
    message  STRING
) WITH (
    'connector' = 'upsert-kafka',
    ...
);

SELECT * FROM upsert_kafka;

이 테이블 소스는 INSERT, UPDATE_AFTER, DELETE 유형의 메시지만 제공하는 반면, 다운스트림 싱크는 완전한 changelog(update_before 포함)를 요구합니다. 결과적으로 이 쿼리 자체는 명시적인 상태 저장 계산을 수반하지 않지만, 플래너는 완전한 changelog를 얻도록 돕기 위해 "ChangelogNormalize"라는 상태 저장 연산자를 여전히 생성합니다.

얼마나 많은 상태가 필요한지, 그리고 잠재적으로 무한히 커지는 상태 크기를 어떻게 제한하는지에 대한 자세한 내용은 개별 연산자 문서를 참고하세요.

유휴 상태 보존 시간 (Idle State Retention Time)

유휴 상태 보존 시간(Idle State Retention Time) 파라미터 table.exec.state.ttl은 키의 상태가 업데이트되지 않고 유지되는 기간을 정의합니다. 위 예시 쿼리에서 word의 개수는 구성된 기간 동안 업데이트되지 않으면 제거됩니다.

키의 상태를 제거하면 연속 쿼리는 이전에 이 키를 본 것을 완전히 잊습니다. 상태가 이전에 제거된 키가 있는 레코드가 처리되면, 그 레코드는 해당 키의 첫 번째 레코드인 것처럼 처리됩니다. 위 예시에서 이는 word의 개수가 다시 0에서 시작함을 의미합니다.

상태 TTL 구성의 다양한 방법 (Different Ways to Configure State TTL)
구성 TableAPI/SQL 지원 세분성 (Granularity) 우선순위 (Priority)
SET 'table.exec.state.ttl' = '...' TableAPI, SQL 파이프라인 레벨, 모든 상태 저장 연산자가 기본적으로 이 값을 사용합니다. 기본 상태 TTL 구성이며, STATE_TTL 힌트를 활성화하거나 직렬화된 CompiledPlan의 값을 수정해 재정의할 수 있습니다.
SELECT /*+ STATE_TTL(...) */ ... SQL 연산자 레벨, 일반 조인과 그룹 집계만 지원합니다. 이 힌트는 기본 table.exec.state.ttl보다 우선합니다. 이 값은 계획 변환 단계에서 CompiledPlan에 직렬화됩니다. State TTL Hint에서 자세히 보세요.
CompiledPlan의 직렬화된 JSON 콘텐츠 수정 TableAPI, SQL 연산자 레벨의 일반화된 지원. 각 상태 저장 연산자의 TTL이 JSON 항목으로 명시적으로 직렬화됩니다. JSON 파일을 수정하면 모든 상태 저장 연산자의 TTL을 변경할 수 있습니다. CompiledPlan의 TTL은 table.exec.state.ttl 또는 STATE_TTL 힌트 중 하나에서 파생됩니다. CompiledPlan으로 작업이 제출되면 최종 TTL 값은 마지막으로 수정된 상태 메타데이터에 의해 결정됩니다.
연산자 레벨 상태 TTL 구성 (Configure Operator-level State TTL)

이것은 고급 기능이므로 주의해서 사용해야 합니다. 파이프라인에서 여러 상태가 사용되고 각 상태에 서로 다른 TTL(Time-to-Live)을 설정해야 하는 경우에만 적합합니다. 파이프라인이 상태 저장 계산을 수반하지 않으면 이 절차를 따를 필요가 없습니다. 파이프라인이 하나의 상태만 사용한다면 파이프라인 레벨에서 table.exec.state.ttl만 설정하면 됩니다.

Table API & SQL은 상태 사용을 개선하기 위해 연산자 레벨의 세밀한 상태 TTL 구성을 지원합니다. 구성 가능한 세분성은 각 상태 연산자의 들어오는 입력 에지 수로 정의됩니다. 구체적으로 OneInputStreamOperator는 하나의 상태에 대한 TTL을 구성할 수 있고, 두 입력이 있는 TwoInputStreamOperator(일반 조인 같은)는 왼쪽과 오른쪽 상태에 대한 TTL을 각각 구성할 수 있습니다. 더 일반적으로 K 입력이 있는 MultipleInputStreamOperator에 대해 K개의 상태 TTL을 구성할 수 있습니다.

전형적인 사용 사례는 다음과 같습니다:

  • 일반 조인에 다른 TTL 설정: 일반 조인은 왼쪽 입력을 유지하는 왼쪽 상태와 오른쪽 입력을 유지하는 오른쪽 상태가 있는 TwoInputStreamOperator를 생성합니다. 왼쪽 상태와 오른쪽 상태에 서로 다른 상태 TTL을 설정할 수 있습니다.
  • 하나의 파이프라인 내 서로 다른 변환에 다른 TTL 설정: 예를 들어 ROW_NUMBER로 중복 제거를 수행한 다음 GROUP BY로 집계를 수행하는 ETL 파이프라인이 있습니다. 이 테이블 프로그램은 각자의 상태를 가진 두 개의 OneInputStreamOperator를 생성합니다. 이제 중복 제거 상태와 집계 상태에 서로 다른 상태 TTL을 설정할 수 있습니다.

윈도우 기반 연산(Window Join, Window Aggregation, Window Top-N 등)과 Interval Join은 상태 보존을 제어하기 위해 table.exec.state.ttl에 의존하지 않으며, 그 상태 TTL은 연산자 레벨에서 구성할 수 없습니다.

Compiled Plan 생성 (Generate a Compiled Plan)

설정 과정은 COMPILE PLAN 문으로 현재 테이블 프로그램의 직렬화된 실행 계획을 나타내는 JSON 파일을 생성하는 것으로 시작합니다.

현재 COMPILE PLAN 문은 SELECT... FROM... 쿼리를 지원하지 않습니다.

  • COMPILE PLAN 문 실행

Java

TableEnvironment tableEnv = TableEnvironment.create(EnvironmentSettings.inStreamingMode());
tableEnv.executeSql(
    "CREATE TABLE orders (order_id BIGINT, order_line_id BIGINT, buyer_id BIGINT, ...)");
tableEnv.executeSql(
    "CREATE TABLE line_orders (order_line_id BIGINT, order_status TINYINT, ...)");
tableEnv.executeSql(
    "CREATE TABLE enriched_orders (order_id BIGINT, order_line_id BIGINT, order_status TINYINT, ...)");

// CompilePlan#writeToFile only supports a local file path, if you need to write to remote filesystem,
// please use tableEnv.executeSql("COMPILE PLAN 'hdfs://path/to/plan.json' FOR ...")
CompiledPlan compiledPlan =
    tableEnv.compilePlanSql(
        "INSERT INTO enriched_orders \n"
       + "SELECT a.order_id, a.order_line_id, b.order_status, ... \n"
       + "FROM orders a JOIN line_orders b ON a.order_line_id = b.order_line_id");

compiledPlan.writeToFile("/path/to/plan.json");

Scala

val tableEnv = TableEnvironment.create(EnvironmentSettings.inStreamingMode())
tableEnv.executeSql(
    "CREATE TABLE orders (order_id BIGINT, order_line_id BIGINT, buyer_id BIGINT, ...)")
tableEnv.executeSql(
    "CREATE TABLE line_orders (order_line_id BIGINT, order_status TINYINT, ...)")
tableEnv.executeSql(
    "CREATE TABLE enriched_orders (order_id BIGINT, order_line_id BIGINT, order_status TINYINT, ...)")

val compiledPlan =
    tableEnv.compilePlanSql(
       """
        |INSERT INTO enriched_orders
        |SELECT a.order_id, a.order_line_id, b.order_status, ...
        |FROM orders a JOIN line_orders b ON a.order_line_id = b.order_line_id
        |""".stripMargin)
// CompilePlan#writeToFile only supports a local file path, if you need to write to remote filesystem,
// please use tableEnv.executeSql("COMPILE PLAN 'hdfs://path/to/plan.json' FOR ...")
compiledPlan.writeToFile("/path/to/plan.json")

SQL CLI

Flink SQL> CREATE TABLE orders (order_id BIGINT, order_line_id BIGINT, buyer_id BIGINT, ...);
[INFO] Execute statement succeeded.

Flink SQL> CREATE TABLE line_orders (order_line_id BIGINT, order_status TINYINT, ...);
[INFO] Execute statement succeeded.

Flink SQL> CREATE TABLE enriched_orders (order_id BIGINT, order_line_id BIGINT, order_status TINYINT, ...);
[INFO] Execute statement succeeded.

Flink SQL> COMPILE PLAN 'file:///path/to/plan.json' FOR INSERT INTO enriched_orders
> SELECT a.order_id, a.order_line_id, b.order_status, ...
> FROM orders a JOIN line_orders b ON a.order_line_id = b.order_line_id;
[INFO] Execute statement succeeded.

SQL 문법 (SQL Syntax)

COMPILE PLAN [IF NOT EXISTS] <plan_file_path> FOR <insert_statement>|<statement_set>;

statement_set:
    EXECUTE STATEMENT SET
    BEGIN
    insert_statement;
    ...
    insert_statement;
    END;

insert_statement:
    <insert_from_select>|<insert_from_values>

이것은 /path/to/plan.json에 JSON 파일을 생성합니다.

COMPILE PLAN 문은 hdfs://s3:// 같은 원격 파일 시스템 스킴으로 계획을 쓰는 것을 지원합니다. 대상 경로에 쓰기 접근이 설정되어 있는지 확인하세요.

Compiled Plan 수정 (Modify the Compiled Plan)

상태를 사용하는 모든 연산자는 다음 구조의 "state"라는 JSON 배열을 명시적으로 생성합니다. 이론적으로 k번째 입력 스트림 연산자는 k번째 상태를 가집니다.

"state": [
    {
      "index": 0,
      "ttl": "0 ms",
      "name": "${1st input state name}"
    },
    {
      "index": 1,
      "ttl": "0 ms",
      "name": "${2nd input state name}"
    },
    ...
  ]

수정해야 할 연산자를 찾아 TTL 값을 양의 정수로 변경하고 시간 단위 "ms"를 포함하는 데 주의하세요. 예를 들어 상태에 1시간을 TTL로 설정하려면 다음처럼 JSON을 수정할 수 있습니다:

{
  "index": 0,
  "ttl": "3600000 ms",
  "name": "${1st input state name}"
}

파일을 저장한 후 EXECUTE PLAN 문으로 작업을 제출합니다.

개념적으로 다운스트림 상태 저장 연산자의 TTL은 업스트림 상태 저장 연산자의 TTL보다 크거나 같아야 합니다.

Compiled Plan 실행 (Execute the Compiled Plan)

EXECUTE PLAN 문은 지정된 파일을 현재 테이블 프로그램의 실행 계획으로 역직렬화한 다음 작업을 제출합니다. EXECUTE PLAN 문으로 제출된 작업은 구성 table.exec.state.ttl 대신 파일에서 읽은 상태 TTL을 적용합니다.

  • EXECUTE PLAN 문 실행

Java

TableEnvironment tableEnv = TableEnvironment.create(EnvironmentSettings.inStreamingMode());
tableEnv.executeSql(
    "CREATE TABLE orders (order_id BIGINT, order_line_id BIGINT, buyer_id BIGINT, ...)");
tableEnv.executeSql(
    "CREATE TABLE line_orders (order_line_id BIGINT, order_status TINYINT, ...)");
tableEnv.executeSql(
    "CREATE TABLE enriched_orders (order_id BIGINT, order_line_id BIGINT, order_status TINYINT, ...)");

// PlanReference#fromFile only supports a local file path, if you need to read from remote filesystem,
// please use tableEnv.executeSql("EXECUTE PLAN 'hdfs://path/to/plan.json'").await();
tableEnv.loadPlan(PlanReference.fromFile("/path/to/plan.json")).execute().await();

Scala

val tableEnv = TableEnvironment.create(EnvironmentSettings.inStreamingMode())
tableEnv.executeSql(
    "CREATE TABLE orders (order_id BIGINT, order_line_id BIGINT, buyer_id BIGINT, ...)")
tableEnv.executeSql(
    "CREATE TABLE line_orders (order_line_id BIGINT, order_status TINYINT, ...)")
tableEnv.executeSql(
    "CREATE TABLE enriched_orders (order_id BIGINT, order_line_id BIGINT, order_status TINYINT, ...)")

// PlanReference#fromFile only supports a local file path, if you need to read from remote filesystem,
// please use tableEnv.executeSql("EXECUTE PLAN 'hdfs://path/to/plan.json'").await()
tableEnv.loadPlan(PlanReference.fromFile("/path/to/plan.json")).execute().await()

SQL CLI

Flink SQL> CREATE TABLE orders (order_id BIGINT, order_line_id BIGINT, buyer_id BIGINT, ...);
[INFO] Execute statement succeeded.

Flink SQL> CREATE TABLE line_orders (order_line_id BIGINT, order_status TINYINT, ...);
[INFO] Execute statement succeeded.

Flink SQL> CREATE TABLE enriched_orders (order_id BIGINT, order_line_id BIGINT, order_status TINYINT, ...);
[INFO] Execute statement succeeded.

Flink SQL> EXECUTE PLAN 'file:///path/to/plan.json';
[INFO] Submitting SQL update statement to the cluster...
[INFO] SQL update statement has been successfully submitted to the cluster:
Job ID: 79fbe3fa497e4689165dd81b1d225ea8

SQL 문법 (SQL Syntax)

EXECUTE PLAN [IF EXISTS] <plan_file_path>;

이것은 JSON 파일을 역직렬화하고 insert statement 작업을 제출합니다.

전체 예시 (A Full Example)

다음 테이블 프로그램은 보강된 주문 출하 정보를 계산합니다. 왼쪽과 오른쪽에 서로 다른 상태 TTL로 일반 내부 조인을 수행합니다.

  • 컴파일된 계획 생성 (Generate compiled plan)
-- left source table
CREATE TABLE Orders (
    `order_id` INT,
    `line_order_id` INT
) WITH (
    'connector'='...'
);

-- right source table
CREATE TABLE LineOrders (
    `line_order_id` INT,
    `ship_mode` STRING
) WITH (
    'connector'='...'
);

-- sink table
CREATE TABLE OrdersShipInfo (
    `order_id` INT,
    `line_order_id` INT,
    `ship_mode` STRING
) WITH (
    'connector' = '...'
);

COMPILE PLAN '/path/to/plan.json' FOR
INSERT INTO OrdersShipInfo
SELECT a.order_id, a.line_order_id, b.ship_mode
FROM Orders a JOIN LineOrders b
    ON a.line_order_id = b.line_order_id;

생성된 JSON 파일은 다음 내용을 가집니다:

{
  "flinkVersion" : "1.18",
  "nodes" : [ {
    "id" : 1,
    "type" : "stream-exec-table-source-scan_1",
    "scanTableSource" : {
      "table" : {
        "identifier" : "`default_catalog`.`default_database`.`Orders`",
        "resolvedTable" : { ... }
      }
    },
    "outputType" : "ROW<`order_id` INT, `line_order_id` INT>",
    "description" : "TableSourceScan(table=[[default_catalog, default_database, Orders]], fields=[order_id, line_order_id])",
    "inputProperties" : [ ]
  }, {
    "id" : 2,
    "type" : "stream-exec-exchange_1",
    "inputProperties" : [ ... ],
    "outputType" : "ROW<`order_id` INT, `line_order_id` INT>",
    "description" : "Exchange(distribution=[hash[line_order_id]])"
  }, {
    "id" : 3,
    "type" : "stream-exec-table-source-scan_1",
    "scanTableSource" : {
      "table" : {
        "identifier" : "`default_catalog`.`default_database`.`LineOrders`",
        "resolvedTable" : {...}
      }
    },
    "outputType" : "ROW<`line_order_id` INT, `ship_mode` VARCHAR(2147483647)>",
    "description" : "TableSourceScan(table=[[default_catalog, default_database, LineOrders]], fields=[line_order_id, ship_mode])",
    "inputProperties" : [ ]
  }, {
    "id" : 4,
    "type" : "stream-exec-exchange_1",
    "inputProperties" : [ ... ],
    "outputType" : "ROW<`line_order_id` INT, `ship_mode` VARCHAR(2147483647)>",
    "description" : "Exchange(distribution=[hash[line_order_id]])"
  }, {
    "id" : 5,
    "type" : "stream-exec-join_1",
    "joinSpec" : { ... },
    "state" : [ {
      "index" : 0,
      "ttl" : "0 ms",
      "name" : "leftState"
    }, {
      "index" : 1,
      "ttl" : "0 ms",
      "name" : "rightState"
    } ],
    "inputProperties" : [ ... ],
    "outputType" : "ROW<`order_id` INT, `line_order_id` INT, `line_order_id0` INT, `ship_mode` VARCHAR(2147483647)>",
    "description" : "Join(joinType=[InnerJoin], where=[(line_order_id = line_order_id0)], select=[order_id, line_order_id, line_order_id0, ship_mode], leftInputSpec=[NoUniqueKey], rightInputSpec=[NoUniqueKey])"
  }, {
    "id" : 6,
    "type" : "stream-exec-calc_1",
    "projection" : [ ... ],
    "condition" : null,
    "inputProperties" : [ ... ],
    "outputType" : "ROW<`order_id` INT, `line_order_id` INT, `ship_mode` VARCHAR(2147483647)>",
    "description" : "Calc(select=[order_id, line_order_id, ship_mode])"
  }, {
    "id" : 7,
    "type" : "stream-exec-sink_1",
    "configuration" : { ... },
    "dynamicTableSink" : {
      "table" : {
        "identifier" : "`default_catalog`.`default_database`.`OrdersShipInfo`",
        "resolvedTable" : { ... }
      }
    },
    "inputChangelogMode" : [ "INSERT" ],
    "inputProperties" : [ ... ],
    "outputType" : "ROW<`order_id` INT, `line_order_id` INT, `ship_mode` VARCHAR(2147483647)>",
    "description" : "Sink(table=[default_catalog.default_database.OrdersShipInfo], fields=[order_id, line_order_id, ship_mode])"
  } ],
  "edges" : [ ... ]
}
  • 계획 콘텐츠 수정 및 계획 실행 (Modify the plan content and execute plan)

조인 연산자 상태의 JSON 표현은 다음 구조를 가집니다:

"state": [
    {
      "index": 0,
      "ttl": "0 ms",
      "name": "leftState"
    },
    {
      "index": 1,
      "ttl": "0 ms",
      "name": "rightState"
    }
  ]

"index"는 현재 상태가 연산자의 i번째 입력임을 나타내며 인덱스는 0부터 시작합니다. 왼쪽과 오른쪽 모두의 현재 TTL 값은 "0 ms"이며, 이는 상태가 만료되지 않음을 의미합니다. 이제 왼쪽 상태 값을 "3000 ms"로, 오른쪽 상태 값을 "9000 ms"로 변경합니다.

"state": [
    {
      "index": 0,
      "ttl": "3000 ms",
      "name": "leftState"
    },
    {
      "index": 1,
      "ttl": "9000 ms",
      "name": "rightState"
    }
  ]

파일에 가한 변경을 저장한 다음 계획을 실행합니다.

EXECUTE PLAN '/path/to/plan.json'

상태 저장 업그레이드와 진화 (Stateful Upgrades and Evolution)

스트리밍 모드에서 실행되는 테이블 프로그램은 상주 쿼리(standing query)로 의도됩니다. 즉, 한 번 정의되고 정적 end-to-end 파이프라인으로 계속 평가됩니다.

상태 저장 파이프라인의 경우 쿼리나 Flink의 플래너에 대한 어떤 변경도 완전히 다른 실행 계획으로 이어질 수 있습니다. 이는 현재 상태 저장 업그레이드와 테이블 프로그램의 진화를 어렵게 만듭니다. 커뮤니티는 이러한 단점을 개선하기 위해 노력하고 있습니다.

예를 들어 필터 술어를 추가하면 옵티마이저가 조인을 재정렬하거나 중간 연산자의 스키마를 변경하기로 결정할 수 있습니다. 이는 토폴로지가 변경되거나 연산자 상태 내의 컬럼 레이아웃이 달라져 세이브포인트에서 복원하는 것을 막습니다.

쿼리 구현자는 변경 전후의 최적화된 계획이 호환되는지 보장해야 합니다. 통찰을 얻으려면 SQL에서 EXPLAIN 명령이나 Table API에서 table.explain()을 사용하세요.

새 옵티마이저 규칙이 계속 추가되고 연산자가 더 효율적이고 특화되므로, 더 새로운 Flink 버전으로 업그레이드하는 것도 호환되지 않는 계획으로 이어질 수 있습니다.

현재 프레임워크는 세이브포인트에서 새 테이블 연산자 토폴로지로 상태를 매핑할 수 있음을 보장할 수 없습니다.

즉, 세이브포인트는 쿼리와 Flink 버전이 모두 일정하게 유지되는 경우에만 지원됩니다.

커뮤니티가 패치 버전(예: 1.13.1에서 1.13.2로)에서 최적화된 계획과 연산자 토폴로지를 수정하는 기여를 거부하므로, Table API & SQL 파이프라인을 더 새로운 버그 수정 릴리스로 업그레이드하는 것은 안전해야 합니다. 그러나 메이저-마이너 업그레이드(예: 1.12에서 1.13으로)는 지원되지 않습니다.

두 단점(즉, 수정된 쿼리와 수정된 Flink 버전) 모두에 대해, 실시간 데이터로 전환하기 전에 업데이트된 테이블 프로그램의 상태를 과거 데이터로 다시 "예열(warm up, 즉 초기화)"할 수 있는지 조사할 것을 권장합니다. Flink 커뮤니티는 이 전환을 가능한 한 편리하게 만들기 위해 하이브리드 소스를 작업 중입니다.

다음 단계 (Where to go next?)

  • Dynamic Tables: 동적 테이블의 개념을 설명합니다.
  • Time attributes: 시간 속성과 Table API & SQL에서 시간 속성이 처리되는 방식을 설명합니다.
  • Versioned Tables: Temporal Table 개념을 설명합니다.
  • Joins in Continuous Queries: 연속 쿼리에서 지원되는 서로 다른 조인 유형.
  • Determinism in Continuous Queries: 연속 쿼리에서의 결정성.
  • Query configuration: Table API & SQL 특화 구성 옵션을 나열합니다.

더 알아보기 (Learn more)