스트리밍 개념

스트리밍 개념 (Streaming Concepts)

Flink의 Table APISQL 지원은 배치(batch)와 스트림(stream) 처리를 위한 통합 API예요. 즉 Table API와 SQL 쿼리는 입력이 유한한 배치 입력이든 무한한 스트림 입력이든 동일한 의미론(semantics)을 가져요.

출처: 문서

본문

다음 페이지들은 스트리밍 데이터에 대한 Flink 관계형 API의 개념, 실질적 한계, 스트림 전용 구성 파라미터를 설명해요.

상태 관리 (State Management)

스트리밍 모드로 실행되는 테이블 프로그램은 상태를 가진 스트림 프로세서로서 Flink의 모든 기능을 활용해요.

특히 테이블 프로그램은 상태 크기와 결함 허용(fault tolerance)에 대한 다양한 요구 사항을 처리하기 위해 state backend와 다양한 체크포인팅 옵션으로 구성할 수 있어요. 실행 중인 Table API & SQL 파이프라인의 savepoint를 찍고, 나중에 시점에서 애플리케이션의 상태를 복원하는 것도 가능해요.

상태 사용 (State Usage)

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

개념적으로 소스 테이블은 상태에 완전히 유지되지 않아요. 구현자는 논리적 테이블(즉 동적 테이블 dynamic tables)을 다뤄요. 그 상태 요구 사항은 사용되는 연산에 따라 달라져요.

상태를 가진 연산자 (Stateful Operators)

조인(joins), 집계(aggregations), 중복 제거(deduplication) 같은 상태를 가진 연산을 포함하는 쿼리는 중간 결과를 결함 허용 저장소에 유지해야 하며, 이를 위해 Flink의 상태 추상화가 사용돼요.

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

또 다른 예시는 단어 개수(word count)를 계산하는 다음 쿼리예요.

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 값이 관찰됨에 따라 쿼리의 전체 상태 크기는 지속적으로 커져요.

Explicit-derived stateful op

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

다음 그림은 upsert kafka source를 질의하는 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"라는 상태를 가진 연산자를 생성해요. Implicit-derived stateful op

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

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

Idle State Retention Time 파라미터 table.exec.state.ttl는 키의 상태가 갱신되지 않고 얼마나 오랫동안 유지된 후 제거되는지를 정의해요. 앞선 예시 쿼리에서 어떤 word의 카운트는 설정된 기간 동안 갱신되지 않으면 즉시 제거돼요.

키의 상태를 제거함으로써 연속 쿼리는 이 키를 이전에 봤다는 것을 완전히 잊어요. 상태가 이전에 제거된 키를 가진 레코드가 처리되면, 그 레코드는 해당 키를 가진 첫 번째 레코드인 것처럼 취급돼요. 위 예시에서는 어떤 word의 카운트가 다시 0에서 시작한다는 뜻이에요.

상태 TTL을 구성하는 서로 다른 방법 (Different Ways to Configure State TTL)

구성(Configuration) 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을 구성할 수 있어요.

전형적인 사용 사례는 다음과 같아요:

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

Window Join, Window Aggregation, Window Top-N 같은 윈도우 기반 연산과 Interval Joins는 상태 보존을 제어하기 위해 table.exec.state.ttl에 의존하지 않으며, 이들의 상태 TTL은 연산자 수준에서 구성할 수 없어요.

컴파일된 플랜 생성 (Generate a Compiled Plan)

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

현재 COMPILE PLAN 문은 SELECT... FROM... 쿼리를 지원하지 않아요.

  • COMPILE PLAN 문 실행
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");
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")
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 문법
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:// 같은 원격 filesystem 스킴으로 작성하는 것을 지원해요. 대상 경로에 쓰기 접근이 설정되어 있는지 확인하세요.

컴파일된 플랜 수정 (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"를 포함하는 것에 주의하세요. 예를 들어 상태의 TTL을 1시간으로 설정하려면 다음과 같이 JSON을 수정할 수 있어요:

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

파일을 저장한 다음 EXECUTE PLAN 문으로 작업을 제출해요.

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

컴파일된 플랜 실행 (Execute the Compiled Plan)

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

  • EXECUTE PLAN 문 실행
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();
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()
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 문법
EXECUTE PLAN [IF EXISTS] <plan_file_path>;

이것은 JSON 파일을 역직렬화하고 insert 문 작업을 제출해요.

완전한 예시 (A Full Example)

다음 테이블 프로그램은 강화된 주문 배송 정보를 계산해요. 왼쪽과 오른쪽 측에 대해 서로 다른 상태 TTL을 가진 일반 내부 조인을 수행해요.

  • 컴파일된 플랜 생성
-- 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" : [ ]
  },
  ...
  "edges" : [ ... ]
}
  • 플랜 내용 수정 및 플랜 실행

조인 연산자의 상태에 대한 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 queries를 목적으로 해요. 즉 한 번 정의되고 정적 종단 간 파이프라인으로 지속적으로 평가된다는 뜻이에요.

상태를 가진 파이프라인의 경우, 쿼리 또는 Flink 플래너에 대한 어떤 변경도 완전히 다른 실행 계획으로 이어질 수 있어요. 이 때문에 현재 상태를 가진 업그레이드와 테이블 프로그램의 진화가 어려워요. 커뮤니티는 이런 단점을 개선하기 위해 노력하고 있어요.

예를 들어 필터 조건자를 추가하면 옵티마이저가 조인 순서를 변경하거나 중간 연산자의 스키마를 변경하기로 결정할 수 있어요. 이것은 위상(topology)이 변경되거나 연산자 상태 내 컬럼 레이아웃이 달라져 savepoint에서 복원하는 것을 방해해요.

쿼리 구현자는 변경 전후의 최적화된 플랜이 호환되는지 확인해야 해요. SQL에서 EXPLAIN 명령 또는 Table API에서 table.explain()을 사용해 통찰을 얻을 수 있어요.

새 옵티마이저 규칙이 지속적으로 추가되고 연산자가 더 효율적이고 특화되기 때문에, 더 새로운 Flink 버전으로의 업그레이드도 호환되지 않는 플랜으로 이어질 수 있어요.

현재 프레임워크는 상태가 savepoint에서 새 테이블 연산자 토폴로지로 매핑될 수 있음을 보장하지 못해요.

다시 말해: savepoint는 쿼리와 Flink 버전이 모두 일정하게 유지되는 경우에만 지원돼요.

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

두 가지 단점(수정된 쿼리와 수정된 Flink 버전) 모두에 대해, 실시간 데이터로 전환하기 전에 업데이트된 테이블 프로그램의 상태를 과거 데이터로 다시 "워밍업"(즉 초기화)할 수 있는지 조사하는 것을 권장해요. Flink 커뮤니티는 이 전환을 최대한 편리하게 만들기 위해 hybrid source를 연구하고 있어요.

다음 단계는 어디로? (Where to go next?)

더 알아보기 (Learn more)