시간 속성

시간 속성 (Time Attributes)

Flink는 서로 다른 시간 개념에 기반하여 데이터를 처리할 수 있습니다. 프로세싱 타임(processing time)과 이벤트 타임(event time) 속성은 Table API와 SQL에서 DDL 또는 DataStream-테이블 변환으로 정의되며, 시간 기반 연산(윈도우, 워터마크 등)에 사용됩니다.

출처: 문서

본문

Flink는 서로 다른 시간 개념에 기반하여 데이터를 처리할 수 있습니다.

  • Processing time은 해당 연산을 실행하는 머신의 시스템 시간(epoch time이라고도 하며, 예: Java의 System.currentTimeMillis())을 말합니다.
  • Event time은 각 행에 첨부된 타임스탬프에 기반한 스트리밍 데이터 처리로, 타임스탬프는 이벤트가 언제 발생했는지를 인코딩할 수 있습니다.

Flink의 시간 처리에 대한 더 많은 정보는 event time and watermarks 소개를 참고하세요.

시간 속성 소개 (Introduction to Time Attributes)

시간 속성은 모든 테이블 스키마의 일부가 될 수 있습니다. 이는 CREATE TABLE DDL 또는 DataStream에서 테이블을 만들 때 정의됩니다. 시간 속성이 정의되면 필드로 참조되어 시간 기반 연산에 사용될 수 있습니다. 시간 속성이 수정되지 않고 질의의 한 부분에서 다른 부분으로 단순히 전달되는 한, 그것은 유효한 시간 속성으로 남습니다. 시간 속성은 일반 타임스탬프처럼 동작하며 계산에 접근 가능합니다. 계산에서 사용될 때 시간 속성은 머티리얼라이즈되어 표준 타임스탬프로 동작합니다. 그러나 일반 타임스탬프는 시간 속성 대신 사용되거나 시간 속성으로 변환될 수 없습니다.

이벤트 타임 (Event Time)

이벤트 타임은 테이블 프로그램이 각 레코드의 타임스탬프에 기반한 결과를 생성할 수 있게 하며, 순서가 뒤바뀌거나 늦은(out-of-order or late) 이벤트에도 일관된 결과를 허용합니다. 또한 영구 저장소에서 레코드를 읽을 때 테이블 프로그램 결과의 재생 가능성(replayability)을 보장합니다.

추가로 이벤트 타임은 배치와 스트리밍 환경 모두에서 테이블 프로그램에 통일된 문법을 허용합니다. 스트리밍 환경의 시간 속성은 배치 환경에서 행의 일반 컬럼이 될 수 있습니다.

스트리밍에서 순서가 뒤바뀐 이벤트를 처리하고 제때 도착한 이벤트와 늦은 이벤트를 구분하려면 Flink가 각 행의 타임스탬프를 알아야 하며, 또한 처리 과정이 이벤트 타임에서 얼마나 진행되었는지에 대한 정기적인 표시(소위 워터마크)도 필요합니다.

이벤트 타임 속성은 CREATE 테이블 DDL 또는 DataStream-테이블 변환 중에 정의할 수 있습니다.

DDL에서 정의 (Defining in DDL)

이벤트 타임 속성은 CREATE 테이블 DDL의 WATERMARK 문을 사용하여 정의합니다. 워터마크 문은 기존 이벤트 타임 필드에 워터마크 생성 표현식을 정의하며, 이벤트 타임 필드를 이벤트 타임 속성으로 표시합니다. 워터마크 문과 워터마크 전략에 대한 자세한 내용은 CREATE TABLE DDL을 참고하세요.

Flink는 TIMESTAMP 컬럼과 TIMESTAMP_LTZ 컬럼에 이벤트 타임 속성 정의를 지원합니다. 소스의 타임스탬프 데이터가 연-월-일-시-분-초로 표현된다면, 보통 시간대 정보 없는 문자열 값(예: 2020-04-15 20:13:40.564)이며, 이벤트 타임 속성을 TIMESTAMP 컬럼으로 정의하는 것을 권장합니다:


CREATE TABLE user_actions (
  user_name STRING,
  data STRING,
  user_action_time TIMESTAMP(3),
  -- declare user_action_time as event time attribute and use 5 seconds delayed watermark strategy
  WATERMARK FOR user_action_time AS user_action_time - INTERVAL '5' SECOND
) WITH (
  ...
);

SELECT TUMBLE_START(user_action_time, INTERVAL '10' MINUTE), COUNT(DISTINCT user_name)
FROM user_actions
GROUP BY TUMBLE(user_action_time, INTERVAL '10' MINUTE);

소스의 타임스탬프 데이터가 epoch time으로 표현된다면, 보통 long 값(예: 1618989564564)이며, 이벤트 타임 속성을 TIMESTAMP_LTZ 컬럼으로 정의하는 것을 권장합니다:


CREATE TABLE user_actions (
  user_name STRING,
  data STRING,
  ts BIGINT,
  time_ltz AS TO_TIMESTAMP_LTZ(ts, 3),
  -- declare time_ltz as event time attribute and use 5 seconds delayed watermark strategy
  WATERMARK FOR time_ltz AS time_ltz - INTERVAL '5' SECOND
) WITH (
  ...
);

SELECT TUMBLE_START(time_ltz, INTERVAL '10' MINUTE), COUNT(DISTINCT user_name)
FROM user_actions
GROUP BY TUMBLE(time_ltz, INTERVAL '10' MINUTE);

고급 워터마크 기능 (Advanced watermark features)

이전 버전에서는 워터마크의 많은 고급 기능(예: 워터마크 정렬)이 datastream API로는 사용하기 쉬웠지만 sql에서는 쉽지 않았습니다. 그래서 1.18 버전에서 이러한 기능을 확장하여 sql에서도 사용할 수 있게 했습니다.

경고 — 참고: SupportsWatermarkPushDown 인터페이스(예: kafka, pulsar)를 구현하는 소스 커넥터만 이러한 고급 기능을 사용할 수 있습니다. 소스가 SupportsWatermarkPushDown 인터페이스를 구현하지 않아도 작업이 이러한 파라미터로 구성되면 작업은 정상적으로 실행될 수 있지만, 파라미터는 적용되지 않습니다.

이러한 기능은 모두 동적 테이블 옵션 또는 'OPTIONS' 힌트로 구성할 수 있습니다. 사용자가 동적 테이블 옵션과 'OPTIONS' 힌트 모두에 기능을 구성했다면 'OPTIONS' 힌트의 옵션이 우선합니다. 사용자가 같은 소스 테이블에 여러 곳에서 'OPTIONS' 힌트를 사용하면 첫 번째 힌트가 사용됩니다.

I. 워터마크 방출 전략 구성 (Configure watermark emit strategy)

flink에서 워터마크를 방출하는 두 가지 전략이 있습니다:

  • on-periodic: 워터마크를 주기적으로 방출.
  • on-event: 이벤트마다 워터마크 방출.

DataStream API에서 사용자는 WatermarkGenerator 인터페이스(Writing WatermarkGenerators)를 통해 방출 전략을 선택할 수 있습니다. sql 작업의 경우 워터마크는 기본적으로 주기적으로 방출되며 기본 주기는 200ms로, pipeline.auto-watermark-interval 파라미터로 변경할 수 있습니다. 이벤트마다 워터마크를 방출해야 한다면 소스 테이블에서 다음과 같이 구성할 수 있습니다:

-- configure in table options
CREATE TABLE user_actions (
  ...
  user_action_time TIMESTAMP(3),
  WATERMARK FOR user_action_time AS user_action_time - INTERVAL '5' SECOND
) WITH (
  'scan.watermark.emit.strategy'='on-event',
  ...
)

물론 OPTIONS 힌트도 사용할 수 있습니다:

-- use 'OPTIONS' hint
select ... from source_table /*+ OPTIONS('scan.watermark.emit.strategy'='on-periodic') */
II. 소스 테이블의 idle-timeout 구성 (Configure the idle-timeout of source table)

소스 테이블의 한 split/partition/shard가 한동안 이벤트 데이터를 보내지 않으면 WatermarkGenerator도 새 워터마크를 생성할 새 데이터를 얻지 못한다는 뜻이며, 우리는 이를 idle input 또는 idle source라고 부릅니다. 이 경우 다른 파티션이 여전히 이벤트 데이터를 보내고 있다면 문제가 발생합니다. 다운스트림 연산자의 워터마크는 모든 업스트림 병렬 데이터 소스 워터마크의 최소값으로 계산되는데, 유휴 split/partition/shard는 새 워터마크를 생성하지 않으므로 다운스트림 연산자의 워터마크가 변하지 않기 때문입니다. 그러나 idle timeout이 구성되면 타임아웃 동안 이벤트 데이터가 없으면 split/partition/shard가 유휴로 표시되고, 다운스트림은 새 워터마크를 계산할 때 이 유휴 소스를 무시합니다.

전역 idle timeout은 sql에서 table.exec.source.idle-timeout 파라미터로 정의할 수 있으며 각 소스 테이블에 적용됩니다. 그러나 소스 테이블마다 다른 idle timeout을 설정하고 싶다면 scan.watermark.idle-timeout 파라미터로 소스 테이블에 다음과 같이 구성할 수 있습니다:

-- configure in table options
CREATE TABLE user_actions (
  ...
  user_action_time TIMESTAMP(3),
  WATERMARK FOR user_action_time AS user_action_time - INTERVAL '5' SECOND
) WITH (
  'scan.watermark.idle-timeout'='1min',
  ...
).

또는 OPTIONS 힌트를 사용할 수 있습니다:

-- use 'OPTIONS' hint
select ... from source_table /*+ OPTIONS('scan.watermark.idle-timeout'='1min') */

사용자가 table.exec.source.idle-timeout 파라미터와 scan.watermark.idle-timeout 파라미터 모두로 소스 idle-timeout을 구성했다면 scan.watermark.idle-timeout 파라미터가 우선합니다.

III. 워터마크 정렬 (Watermark alignment)

데이터 분포나 머신 부하 같은 다양한 요인의 영향으로 같은 데이터 소스 또는 다른 데이터 소스의 서로 다른 split/partition/shard 간 소비 속도가 다를 수 있습니다. 다운스트림에 상태 저장 연산자가 있다면 이러한 연산자는 더 빨리 소비하는 것에 대해 상태에 더 많은 데이터를 캐시하고 더 느리게 소비하는 것을 기다려야 하므로 상태가 매우 커질 수 있습니다. 불일치한 소비 속도는 더 심각한 데이터 무질서를 유발해 윈도우 계산 정확도에 영향을 줄 수 있습니다. 워터마크 정렬 기능을 사용하면 빠른 split/partition/shard의 워터마크가 다른 split/partition/shard에 비해 너무 빨리 증가하지 않게 하여 이러한 시나리오를 피할 수 있습니다. 워터마크 정렬 기능은 서로 다른 소스 간 데이터 소비 차이에 따라 작업 성능에 영향을 준다는 점에 유의하세요.

워터마크 정렬은 소스 테이블에서 다음과 같이 구성할 수 있습니다:

-- configure in table options
CREATE TABLE user_actions (
...
user_action_time TIMESTAMP(3),
  WATERMARK FOR user_action_time AS user_action_time - INTERVAL '5' SECOND
) WITH (
'scan.watermark.alignment.group'='alignment-group-1',
'scan.watermark.alignment.max-drift'='1min',
'scan.watermark.alignment.update-interval'='1s',
...
).

물론 여전히 OPTIONS 힌트를 사용할 수 있습니다.

-- use 'OPTIONS' hint
select ... from source_table /*+ OPTIONS('scan.watermark.alignment.group'='alignment-group-1', 'scan.watermark.alignment.max-drift'='1min', 'scan. watermark.alignment.update-interval'='1s') */

세 가지 파라미터가 있습니다:

  • scan.watermark.alignment.group은 정렬 그룹 이름을 구성하며, 같은 그룹의 데이터 소스는 정렬됩니다
  • scan.watermark.alignment.max-drift는 split/partition/shard에 허용되는 정렬 시간에서의 최대 편차 범위를 구성합니다
  • scan.watermark.alignment.update-interval은 정렬 시간이 계산되는 빈도를 구성하며, 필수가 아니고 기본값은 1s입니다

경고 — 참고: FLIP-217에 따라 1.17부터 워터마크 정렬 기능을 사용하려면 커넥터가 source split의 워터마크 정렬을 구현해야 합니다. 소스 커넥터가 FLIP-217을 구현하지 않으면 작업이 오류와 함께 실행됩니다. 사용자는 pipeline.watermark-alignment.allow-unaligned-source-splits: true를 설정하여 source split의 워터마크 정렬을 비활성화할 수 있으며, 워터마크 정렬은 split 수가 소스 연산자의 병렬도와 같을 때만 제대로 동작합니다.

DataStream-테이블 변환 중 (During DataStream-to-Table Conversion)

DataStream을 테이블로 변환할 때 이벤트 타임 속성은 스키마 정의 중 .rowtime 속성으로 정의할 수 있습니다. 타임스탬프와 워터마크는 변환되는 DataStream에 이미 할당되어 있어야 합니다. 변환 중 Flink는 항상 rowtime 속성을 TIMESTAMP WITHOUT TIME ZONE으로 파생합니다. DataStream에는 시간대 개념이 없고 모든 이벤트 타임 값을 UTC로 처리하기 때문입니다.

DataStreamTable로 변환할 때 시간 속성을 정의하는 두 가지 방법이 있습니다. 지정된 .rowtime 필드 이름이 DataStream의 스키마에 존재하는지 여부에 따라 타임스탬프는 (1) 새 컬럼으로 추가되거나 (2) 기존 컬럼을 대체합니다.

어느 경우든 이벤트 타임 타임스탬프 필드는 DataStream 이벤트 타임 타임스탬프의 값을 보유합니다.

Java:


// Option 1:

// extract timestamp and assign watermarks based on knowledge of the stream
DataStream<Tuple2<String, String>> stream = inputStream.assignTimestampsAndWatermarks(...);

// declare an additional logical field as an event time attribute
Table table = tEnv.fromDataStream(stream, $("user_name"), $("data"), $("user_action_time").rowtime());


// Option 2:

// extract timestamp from first field, and assign watermarks based on knowledge of the stream
DataStream<Tuple3<Long, String, String>> stream = inputStream.assignTimestampsAndWatermarks(...);

// the first field has been used for timestamp extraction, and is no longer necessary
// replace first field with a logical event time attribute
Table table = tEnv.fromDataStream(stream, $("user_action_time").rowtime(), $("user_name"), $("data"));

// Usage:

WindowedTable windowedTable = table.window(Tumble
       .over(lit(10).minutes())
       .on($("user_action_time"))
       .as("userActionWindow"));

Scala:


// Option 1:

// extract timestamp and assign watermarks based on knowledge of the stream
val stream: DataStream[(String, String)] = inputStream.assignTimestampsAndWatermarks(...)

// declare an additional logical field as an event time attribute
val table = tEnv.fromDataStream(stream, $"user_name", $"data", $"user_action_time".rowtime)


// Option 2:

// extract timestamp from first field, and assign watermarks based on knowledge of the stream
val stream: DataStream[(Long, String, String)] = inputStream.assignTimestampsAndWatermarks(...)

// the first field has been used for timestamp extraction, and is no longer necessary
// replace first field with a logical event time attribute
val table = tEnv.fromDataStream(stream, $"user_action_time".rowtime, $"user_name", $"data")

// Usage:

val windowedTable = table.window(Tumble over 10.minutes on $"user_action_time" as "userActionWindow")

Python:


# Option 1:

# extract timestamp and assign watermarks based on knowledge of the stream
stream = input_stream.assign_timestamps_and_watermarks(...)

table = t_env.from_data_stream(stream, col('user_name'), col('data'), col('user_action_time').rowtime)

# Option 2:

# extract timestamp from first field, and assign watermarks based on knowledge of the stream
stream = input_stream.assign_timestamps_and_watermarks(...)

# the first field has been used for timestamp extraction, and is no longer necessary
# replace first field with a logical event time attribute
table = t_env.from_data_stream(stream, col("user_action_time").rowtime, col('user_name'), col('data'))

# Usage:

table.window(Tumble.over(lit(10).minutes).on(col("user_action_time")).alias("userActionWindow"))

프로세싱 타임 (Processing Time)

프로세싱 타임은 테이블 프로그램이 로컬 머신의 시간에 기반한 결과를 생성할 수 있게 합니다. 가장 단순한 시간 개념이지만 비결정적(non-deterministic) 결과를 생성합니다. 프로세싱 타임은 타임스탬프 추출이나 워터마크 생성이 필요하지 않습니다.

프로세싱 타임 속성을 정의하는 두 가지 방법이 있습니다.

DDL에서 정의 (Defining in DDL)

프로세싱 타임 속성은 CREATE 테이블 DDL에서 시스템 PROCTIME() 함수를 사용한 계산 컬럼(computed column)으로 정의되며, 함수의 반환 타입은 TIMESTAMP_LTZ입니다. 계산 컬럼에 대한 자세한 내용은 CREATE TABLE DDL을 참고하세요.


CREATE TABLE user_actions (
  user_name STRING,
  data STRING,
  user_action_time AS PROCTIME() -- declare an additional field as a processing time attribute
) WITH (
  ...
);

SELECT TUMBLE_START(user_action_time, INTERVAL '10' MINUTE), COUNT(DISTINCT user_name)
FROM user_actions
GROUP BY TUMBLE(user_action_time, INTERVAL '10' MINUTE);

DataStream-테이블 변환 중 (During DataStream-to-Table Conversion)

프로세싱 타임 속성은 스키마 정의 중 .proctime 속성으로 정의됩니다. 시간 속성은 물리적 스키마를 추가 논리 필드로만 확장해야 합니다. 따라서 스키마 정의의 끝에서만 정의할 수 있습니다.

Java:

DataStream<Tuple2<String, String>> stream = ...;

// declare an additional logical field as a processing time attribute
Table table = tEnv.fromDataStream(stream, $("user_name"), $("data"), $("user_action_time").proctime());

WindowedTable windowedTable = table.window(
        Tumble.over(lit(10).minutes())
            .on($("user_action_time"))
            .as("userActionWindow"));

Scala:

val stream: DataStream[(String, String)] = ...

// declare an additional logical field as a processing time attribute
val table = tEnv.fromDataStream(stream, $"UserActionTimestamp", $"user_name", $"data", $"user_action_time".proctime)

val windowedTable = table.window(Tumble over 10.minutes on $"user_action_time" as "userActionWindow")

Python:

stream = ...

# declare an additional logical field as a processing time attribute
table = t_env.from_data_stream(stream, col("UserActionTimestamp"), col("user_name"), col("data"), col("user_action_time").proctime)

windowed_table = table.window(Tumble.over(lit(10).minutes).on(col("user_action_time")).alias("userActionWindow"))

더 알아보기 (Learn more)