시간 속성

시간 속성 (Time Attributes)

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

  • Processing time은 각각의 연산을 실행하는 머신의 시스템 시간(에포크 시간이라고도 함, 예: Java의 System.currentTimeMillis())을 가리킵니다.
  • Event time은 각 행에 부착된 타임스탬프를 기반으로 스트리밍 데이터를 처리하는 것입니다. 타임스탬프는 이벤트가 발생한 시점을 인코딩할 수 있습니다.

Flink의 시간 처리에 대한 자세한 내용은 event time and watermarks 소개를 참조하세요.

출처: 문서

본문

시간 속성 소개

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

이벤트 시간 (Event Time)

이벤트 시간은 테이블 프로그램이 모든 레코드의 타임스탬프를 기반으로 결과를 생성할 수 있게 하여, 순서가 바뀌거나 지연된 이벤트에도 일관된 결과를 허용합니다. 또한 영구 저장소에서 레코드를 읽을 때 테이블 프로그램 결과의 재생 가능성도 보장합니다.

또한 이벤트 시간은 배치와 스트리밍 환경 모두에서 테이블 프로그램에 대한 통합 구문을 허용합니다. 스트리밍 환경의 시간 속성은 배치 환경에서 행의 일반 열일 수 있습니다.

스트리밍에서 순서가 바뀐 이벤트를 처리하고 제시간(on-time) 및 지연(late) 이벤트를 구분하기 위해 Flink는 각 행의 타임스탬프와 이벤트 시간 처리가 여기까지 얼마나 진행되었는지에 대한 정기적인 표시(소위 워터마크)가 필요합니다.

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

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);

소스의 타임스탬프 데이터가 에포크 시간(보통 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);
고급 워터마크 기능

이전 버전에서는 워터마크의 많은 고급 기능(예: 워터마크 정렬)이 datastream API를 통해서는 사용하기 쉬웠지만 sql에서는 그렇지 않았습니다. 그래서 1.18 버전에서 사용자가 sql에서도 사용할 수 있도록 이 기능들을 확장했습니다.

참고: SupportsWatermarkPushDown 인터페이스를 구현하는 소스 커넥터(예: kafka, pulsar)만 이 고급 기능을 사용할 수 있습니다. 소스가 SupportsWatermarkPushDown 인터페이스를 구현하지 않지만 작업에 이 매개변수가 구성된 경우 작업은 정상적으로 실행될 수 있지만 이 매개변수는 효과가 없습니다.

이 기능들은 모두 dynamic table options 또는 'OPTIONS' 힌트로 구성할 수 있습니다. 사용자가 dynamic table options와 'OPTIONS' 힌트 양쪽 모두에 이 기능을 구성했다면 'OPTIONS' 힌트의 옵션이 우선합니다. 사용자가 같은 소스 테이블에 대해 여러 곳에서 'OPTIONS' 힌트를 사용하면 첫 번째 힌트가 사용됩니다.

I. 워터마크 방출 전략 구성

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 구성

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

전역 idle 타임아웃은 sql에서 table.exec.source.idle-timeout 매개변수로 정의할 수 있으며, 이는 각 소스 테이블에 적용됩니다. 그러나 소스 테이블마다 다른 idle 타임아웃을 설정하려면 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 사이에 소비 속도가 다를 수 있습니다. 다운스트림에 상태 기반 연산자가 있다면 이 연산자는 더 빨리 소비하는 쪽에 대해 상태에 더 많은 데이터를 캐시하고 더 느리게 소비하는 쪽을 기다려야 할 수 있으므로 상태가 매우 커질 수 있습니다. 일관되지 않은 소비 속도는 더 심각한 데이터 무질서(order)를 유발해 윈도우의 계산 정확성에 영향을 줄 수 있습니다. 워터마크 정렬 기능을 사용하면 빠른 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부터 소스 split의 워터마크 정렬을 사용하려면 커넥터가 이를 구현해야 합니다. 소스 커넥터가 FLIP-217을 구현하지 않으면 작업은 오류와 함께 실행되며, 사용자는 pipeline.watermark-alignment.allow-unaligned-source-splits: true를 설정하여 소스 split의 워터마크 정렬을 비활성화할 수 있습니다. 워터마크 정렬은 split 수가 소스 연산자의 병렬도와 같을 때만 제대로 작동합니다.

DataStream-to-Table 변환 중

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

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)

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

처리 시간 속성을 정의하는 방법은 두 가지가 있습니다.

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-to-Table 변환 중

처리 시간 속성은 스키마 정의 중 .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)