시간대
시간대 (Time Zone)
Flink는 DATE, TIME, TIMESTAMP, TIMESTAMP_LTZ, INTERVAL YEAR TO MONTH, INTERVAL DAY TO SECOND 같은 풍부한 날짜/시간 데이터 타입을 제공해요 (자세한 내용은 Date and Time 참고).
Flink는 세션 수준에서 시간대를 설정하는 것을 지원해요 (자세한 내용은 table.local-time-zone 참고).
이런 타임스탬프 데이터 타입과 Flink의 시간대 지원 덕분에 시간대를 넘나드는 비즈니스 데이터를 쉽게 처리할 수 있어요.
출처: 문서
본문
TIMESTAMP vs TIMESTAMP_LTZ
TIMESTAMP 타입
TIMESTAMP(p)는TIMESTAMP(p) WITHOUT TIME ZONE의 약어예요. 정밀도p는 0에서 9까지 범위를 지원하며, 기본값은 6이에요.TIMESTAMP는 연, 월, 일, 시, 분, 초 및 분수 초(fractional seconds)를 나타내는 타임스탬프를 설명해요.TIMESTAMP는 문자열 리터럴에서 지정될 수 있어요. 예:
Flink SQL> SELECT TIMESTAMP '1970-01-01 00:00:04.001';
+-------------------------+
| 1970-01-01 00:00:04.001 |
+-------------------------+
TIMESTAMP_LTZ 타입
TIMESTAMP_LTZ(p)는TIMESTAMP(p) WITH LOCAL TIME ZONE의 약어예요. 정밀도p는 0에서 9까지 범위를 지원하며, 기본값은 6이에요.TIMESTAMP_LTZ는 시간선(time-line)에서의 절대 시점(absolute time point)을 설명해요. epoch-밀리초를 나타내는 long 값과 밀리초의 나노초를 나타내는 int 값을 저장해요. epoch 시간은 표준 Java epoch인1970-01-01T00:00:00Z에서 측정돼요.TIMESTAMP_LTZ타입의 모든 데이터는 계산과 시각화를 위해 현재 세션에 구성된 로컬 시간대로 해석돼요.TIMESTAMP_LTZ는 리터럴 표현이 없으므로 리터럴로부터 지정할 수 없어요. long epoch 시간(예: JavaSystem.currentTimeMillis()가 생성한 long 시간)에서 파생될 수 있어요.
Flink SQL> CREATE VIEW T1 AS SELECT TO_TIMESTAMP_LTZ(4001, 3);
Flink SQL> SET 'table.local-time-zone' = 'UTC';
Flink SQL> SELECT * FROM T1;
+---------------------------+
| TO_TIMESTAMP_LTZ(4001, 3) |
+---------------------------+
| 1970-01-01 00:00:04.001 |
+---------------------------+
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT * FROM T1;
+---------------------------+
| TO_TIMESTAMP_LTZ(4001, 3) |
+---------------------------+
| 1970-01-01 08:00:04.001 |
+---------------------------+
TIMESTAMP_LTZ는 절대 시점(예: 위의4001밀리초)이 서로 다른 시간대에서 동일한 순간(instantaneous point)을 설명하므로, 시간대를 넘나드는 비즈니스에 사용될 수 있어요. 같은 시점에 세계 모든 머신의System.currentTimeMillis()가 같은 값을 반환한다는 배경(예: 위 예시의4001밀리초)을 상기하면, 이것이 절대 시점 의미라는 것을 알 수 있어요.
시간대 사용 (Time Zone Usage)
로컬 시간대는 현재 세션 시간대 id를 정의해요. SQL Client 또는 애플리케이션에서 시간대를 구성할 수 있어요.
-- set to UTC time zone
Flink SQL> SET 'table.local-time-zone' = 'UTC';
-- set to Shanghai time zone
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
-- set to Los_Angeles time zone
Flink SQL> SET 'table.local-time-zone' = 'America/Los_Angeles';
EnvironmentSettings envSetting = EnvironmentSettings.inStreamingMode();
TableEnvironment tEnv = TableEnvironment.create(envSetting);
// set to UTC time zone
tEnv.getConfig().setLocalTimeZone(ZoneId.of("UTC"));
// set to Shanghai time zone
tEnv.getConfig().setLocalTimeZone(ZoneId.of("Asia/Shanghai"));
// set to Los_Angeles time zone
tEnv.getConfig().setLocalTimeZone(ZoneId.of("America/Los_Angeles"));
val envSetting = EnvironmentSettings.inStreamingMode()
val tEnv = TableEnvironment.create(envSetting)
// set to UTC time zone
tEnv.getConfig.setLocalTimeZone(ZoneId.of("UTC"))
// set to Shanghai time zone
tEnv.getConfig.setLocalTimeZone(ZoneId.of("Asia/Shanghai"))
// set to Los_Angeles time zone
tEnv.getConfig.setLocalTimeZone(ZoneId.of("America/Los_Angeles"))
env_setting = EnvironmentSettings.in_streaming_mode()
t_env = TableEnvironment.create(env_setting)
# set to UTC time zone
t_env.get_config().set_local_timezone("UTC")
# set to Shanghai time zone
t_env.get_config().set_local_timezone("Asia/Shanghai")
# set to Los_Angeles time zone
t_env.get_config().set_local_timezone("America/Los_Angeles")
세션 시간대는 Flink SQL에서 유용한데, 주요 용도는 다음과 같아요:
시간 함수 반환 값 결정 (Decide time functions return value)
다음 시간 함수들은 구성된 시간대의 영향을 받아요:
- LOCALTIME
- LOCALTIMESTAMP
- CURRENT_DATE
- CURRENT_TIME
- CURRENT_TIMESTAMP
- CURRENT_ROW_TIMESTAMP()
- NOW()
- PROCTIME()
Flink SQL> SET 'sql-client.execution.result-mode' = 'tableau';
Flink SQL> CREATE VIEW MyView1 AS SELECT LOCALTIME, LOCALTIMESTAMP, CURRENT_DATE, CURRENT_TIME, CURRENT_TIMESTAMP, CURRENT_ROW_TIMESTAMP(), NOW(), PROCTIME();
Flink SQL> DESC MyView1;
+------------------------+-----------------------------+-------+-----+--------+-----------+
| name | type | null | key | extras | watermark |
+------------------------+-----------------------------+-------+-----+--------+-----------+
| LOCALTIME | TIME(0) | false | | | |
| LOCALTIMESTAMP | TIMESTAMP(3) | false | | | |
| CURRENT_DATE | DATE | false | | | |
| CURRENT_TIME | TIME(0) | false | | | |
| CURRENT_TIMESTAMP | TIMESTAMP_LTZ(3) | false | | | |
|CURRENT_ROW_TIMESTAMP() | TIMESTAMP_LTZ(3) | false | | | |
| NOW() | TIMESTAMP_LTZ(3) | false | | | |
| PROCTIME() | TIMESTAMP_LTZ(3) *PROCTIME* | false | | | |
+------------------------+-----------------------------+-------+-----+--------+-----------+
Flink SQL> SET 'table.local-time-zone' = 'UTC';
Flink SQL> SELECT * FROM MyView1;
+-----------+-------------------------+--------------+--------------+-------------------------+-------------------------+-------------------------+-------------------------+
| LOCALTIME | LOCALTIMESTAMP | CURRENT_DATE | CURRENT_TIME | CURRENT_TIMESTAMP | CURRENT_ROW_TIMESTAMP() | NOW() | PROCTIME() |
+-----------+-------------------------+--------------+--------------+-------------------------+-------------------------+-------------------------+-------------------------+
| 15:18:36 | 2021-04-15 15:18:36.384 | 2021-04-15 | 15:18:36 | 2021-04-15 15:18:36.384 | 2021-04-15 15:18:36.384 | 2021-04-15 15:18:36.384 | 2021-04-15 15:18:36.384 |
+-----------+-------------------------+--------------+--------------+-------------------------+-------------------------+-------------------------+-------------------------+
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT * FROM MyView1;
+-----------+-------------------------+--------------+--------------+-------------------------+-------------------------+-------------------------+-------------------------+
| LOCALTIME | LOCALTIMESTAMP | CURRENT_DATE | CURRENT_TIME | CURRENT_TIMESTAMP | CURRENT_ROW_TIMESTAMP() | NOW() | PROCTIME() |
+-----------+-------------------------+--------------+--------------+-------------------------+-------------------------+-------------------------+-------------------------+
| 23:18:36 | 2021-04-15 23:18:36.384 | 2021-04-15 | 23:18:36 | 2021-04-15 23:18:36.384 | 2021-04-15 23:18:36.384 | 2021-04-15 23:18:36.384 | 2021-04-15 23:18:36.384 |
+-----------+-------------------------+--------------+--------------+-------------------------+-------------------------+-------------------------+-------------------------+
TIMESTAMP_LTZ 문자열 표현
TIMESTAMP_LTZ 값을 문자열 형식으로 표현할 때(즉, 값을 출력하거나, STRING 타입으로 캐스트하거나, TIMESTAMP로 캐스트하거나, TIMESTAMP 값을 TIMESTAMP_LTZ로 캐스트할 때) 세션 시간대가 사용돼요:
Flink SQL> CREATE VIEW MyView2 AS SELECT TO_TIMESTAMP_LTZ(4001, 3) AS ltz, TIMESTAMP '1970-01-01 00:00:01.001' AS ntz;
Flink SQL> DESC MyView2;
+------+------------------+-------+-----+--------+-----------+
| name | type | null | key | extras | watermark |
+------+------------------+-------+-----+--------+-----------+
| ltz | TIMESTAMP_LTZ(3) | true | | | |
| ntz | TIMESTAMP(3) | false | | | |
+------+------------------+-------+-----+--------+-----------+
Flink SQL> SET 'table.local-time-zone' = 'UTC';
Flink SQL> SELECT * FROM MyView2;
+-------------------------+-------------------------+
| ltz | ntz |
+-------------------------+-------------------------+
| 1970-01-01 00:00:04.001 | 1970-01-01 00:00:01.001 |
+-------------------------+-------------------------+
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT * FROM MyView2;
+-------------------------+-------------------------+
| ltz | ntz |
+-------------------------+-------------------------+
| 1970-01-01 08:00:04.001 | 1970-01-01 00:00:01.001 |
+-------------------------+-------------------------+
Flink SQL> CREATE VIEW MyView3 AS SELECT ltz, CAST(ltz AS TIMESTAMP(3)), CAST(ltz AS STRING), ntz, CAST(ntz AS TIMESTAMP_LTZ(3)) FROM MyView2;
Flink SQL> DESC MyView3;
+-------------------------------+------------------+-------+-----+--------+-----------+
| name | type | null | key | extras | watermark |
+-------------------------------+------------------+-------+-----+--------+-----------+
| ltz | TIMESTAMP_LTZ(3) | true | | | |
| CAST(ltz AS TIMESTAMP(3)) | TIMESTAMP(3) | true | | | |
| CAST(ltz AS STRING) | STRING | true | | | |
| ntz | TIMESTAMP(3) | false | | | |
| CAST(ntz AS TIMESTAMP_LTZ(3)) | TIMESTAMP_LTZ(3) | false | | | |
+-------------------------------+------------------+-------+-----+--------+-----------+
Flink SQL> SELECT * FROM MyView3;
+-------------------------+---------------------------+-------------------------+-------------------------+-------------------------------+
| ltz | CAST(ltz AS TIMESTAMP(3)) | CAST(ltz AS STRING) | ntz | CAST(ntz AS TIMESTAMP_LTZ(3)) |
+-------------------------+---------------------------+-------------------------+-------------------------+-------------------------------+
| 1970-01-01 08:00:04.001 | 1970-01-01 08:00:04.001 | 1970-01-01 08:00:04.001 | 1970-01-01 00:00:01.001 | 1970-01-01 00:00:01.001 |
+-------------------------+---------------------------+-------------------------+-------------------------+-------------------------------+
시간 속성과 시간대 (Time Attribute and Time Zone)
시간 속성에 대한 자세한 내용은 Time Attribute를 참조하세요.
프로세싱 타임과 시간대 (Processing Time and Time Zone)
Flink SQL은 PROCTIME() 함수로 프로세스 타임 속성을 정의하며, 함수 반환 타입은 TIMESTAMP_LTZ예요.
Flink 1.13 이전에는
PROCTIME()의 함수 반환 타입이TIMESTAMP이고, 반환 값이 UTC 시간대의TIMESTAMP였어요. 예를 들어 상하이에서 벽시계가2021-03-01 12:00:00을 가리키는데PROCTIME()이2021-03-01 04:00:00을 표시하는 것은 잘못된 것이에요. Flink 1.13은 이 문제를 수정하고PROCTIME()의 반환 타입으로TIMESTAMP_LTZ를 사용하므로, 사용자는 더 이상 시간대 문제를 처리할 필요가 없어요.
PROCTIME()은 항상 로컬 타임스탬프 값을 나타내며, TIMESTAMP_LTZ 타입을 사용하면 일광 절약 시간(Daylight Saving Time)도 잘 지원할 수 있어요.
Flink SQL> SET 'table.local-time-zone' = 'UTC';
Flink SQL> SELECT PROCTIME();
+-------------------------+
| PROCTIME() |
+-------------------------+
| 2021-04-15 14:48:31.387 |
+-------------------------+
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT PROCTIME();
+-------------------------+
| PROCTIME() |
+-------------------------+
| 2021-04-15 22:48:31.387 |
+-------------------------+
Flink SQL> CREATE TABLE MyTable1 (
item STRING,
price DOUBLE,
proctime as PROCTIME()
) WITH (
'connector' = 'socket',
'hostname' = '127.0.0.1',
'port' = '9999',
'format' = 'csv'
);
Flink SQL> CREATE VIEW MyView3 AS
SELECT
TUMBLE_START(proctime, INTERVAL '10' MINUTES) AS window_start,
TUMBLE_END(proctime, INTERVAL '10' MINUTES) AS window_end,
TUMBLE_PROCTIME(proctime, INTERVAL '10' MINUTES) as window_proctime,
item,
MAX(price) as max_price
FROM MyTable1
GROUP BY TUMBLE(proctime, INTERVAL '10' MINUTES), item;
Flink SQL> DESC MyView3;
+-----------------+-----------------------------+-------+-----+--------+-----------+
| name | type | null | key | extras | watermark |
+-----------------+-----------------------------+-------+-----+--------+-----------+
| window_start | TIMESTAMP(3) | false | | | |
| window_end | TIMESTAMP(3) | false | | | |
| window_proctime | TIMESTAMP_LTZ(3) *PROCTIME* | false | | | |
| item | STRING | true | | | |
| max_price | DOUBLE | true | | | |
+-----------------+-----------------------------+-------+-----+--------+-----------+
터미널에서 다음 명령을 사용해 MyTable1의 데이터를 수집해요:
> nc -lk 9999
A,1.1
B,1.2
A,1.8
B,2.5
C,3.8
Flink SQL> SET 'table.local-time-zone' = 'UTC';
Flink SQL> SELECT * FROM MyView3;
+-------------------------+-------------------------+-------------------------+------+-----------+
| window_start | window_end | window_procime | item | max_price |
+-------------------------+-------------------------+-------------------------+------+-----------+
| 2021-04-15 14:00:00.000 | 2021-04-15 14:10:00.000 | 2021-04-15 14:10:00.005 | A | 1.8 |
| 2021-04-15 14:00:00.000 | 2021-04-15 14:10:00.000 | 2021-04-15 14:10:00.007 | B | 2.5 |
| 2021-04-15 14:00:00.000 | 2021-04-15 14:10:00.000 | 2021-04-15 14:10:00.007 | C | 3.8 |
+-------------------------+-------------------------+-------------------------+------+-----------+
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT * FROM MyView3;
UTC 시간대에서의 계산과 비교해 다른 window start, window end, window proctime을 반환해요.
+-------------------------+-------------------------+-------------------------+------+-----------+
| window_start | window_end | window_procime | item | max_price |
+-------------------------+-------------------------+-------------------------+------+-----------+
| 2021-04-15 22:00:00.000 | 2021-04-15 22:10:00.000 | 2021-04-15 22:10:00.005 | A | 1.8 |
| 2021-04-15 22:00:00.000 | 2021-04-15 22:10:00.000 | 2021-04-15 22:10:00.007 | B | 2.5 |
| 2021-04-15 22:00:00.000 | 2021-04-15 22:10:00.000 | 2021-04-15 22:10:00.007 | C | 3.8 |
+-------------------------+-------------------------+-------------------------+------+-----------+
프로세싱 타임 윈도우는 비결정적(non-deterministic)이므로, 실행할 때마다 다른 윈도우와 다른 집계를 얻을 거예요. 위 예시는 시간대가 프로세싱 타임 윈도우에 어떻게 영향을 주는지 설명하기 위한 것이에요.
이벤트 타임과 시간대 (Event Time and Time Zone)
Flink는 TIMESTAMP 컬럼과 TIMESTAMP_LTZ 컬럼에 이벤트 타임 속성을 정의하는 것을 지원해요.
TIMESTAMP의 이벤트 타임 속성
소스의 타임스탬프 데이터가 연-월-일-시-분-초로 표현된다면(보통 시간대 정보가 없는 문자열 값, 예: 2020-04-15 20:13:40.564), 이벤트 타임 속성을 TIMESTAMP 컬럼으로 정의하는 것을 권장해요:
Flink SQL> CREATE TABLE MyTable2 (
item STRING,
price DOUBLE,
ts TIMESTAMP(3), -- TIMESTAMP data type
WATERMARK FOR ts AS ts - INTERVAL '10' SECOND
) WITH (
'connector' = 'socket',
'hostname' = '127.0.0.1',
'port' = '9999',
'format' = 'csv'
);
Flink SQL> CREATE VIEW MyView4 AS
SELECT
TUMBLE_START(ts, INTERVAL '10' MINUTES) AS window_start,
TUMBLE_END(ts, INTERVAL '10' MINUTES) AS window_end,
TUMBLE_ROWTIME(ts, INTERVAL '10' MINUTES) as window_rowtime,
item,
MAX(price) as max_price
FROM MyTable2
GROUP BY TUMBLE(ts, INTERVAL '10' MINUTES), item;
Flink SQL> DESC MyView4;
+----------------+------------------------+------+-----+--------+-----------+
| name | type | null | key | extras | watermark |
+----------------+------------------------+------+-----+--------+-----------+
| window_start | TIMESTAMP(3) | true | | | |
| window_end | TIMESTAMP(3) | true | | | |
| window_rowtime | TIMESTAMP(3) *ROWTIME* | true | | | |
| item | STRING | true | | | |
| max_price | DOUBLE | true | | | |
+----------------+------------------------+------+-----+--------+-----------+
터미널에서 다음 명령을 사용해 MyTable2의 데이터를 수집해요:
> nc -lk 9999
A,1.1,2021-04-15 14:01:00
B,1.2,2021-04-15 14:02:00
A,1.8,2021-04-15 14:03:00
B,2.5,2021-04-15 14:04:00
C,3.8,2021-04-15 14:05:00
C,3.8,2021-04-15 14:11:00
Flink SQL> SET 'table.local-time-zone' = 'UTC';
Flink SQL> SELECT * FROM MyView4;
+-------------------------+-------------------------+-------------------------+------+-----------+
| window_start | window_end | window_rowtime | item | max_price |
+-------------------------+-------------------------+-------------------------+------+-----------+
| 2021-04-15 14:00:00.000 | 2021-04-15 14:10:00.000 | 2021-04-15 14:09:59.999 | A | 1.8 |
| 2021-04-15 14:00:00.000 | 2021-04-15 14:10:00.000 | 2021-04-15 14:09:59.999 | B | 2.5 |
| 2021-04-15 14:00:00.000 | 2021-04-15 14:10:00.000 | 2021-04-15 14:09:59.999 | C | 3.8 |
+-------------------------+-------------------------+-------------------------+------+-----------+
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT * FROM MyView4;
UTC 시간대에서의 계산과 비교해 동일한 window start, window end, window rowtime을 반환해요.
+-------------------------+-------------------------+-------------------------+------+-----------+
| window_start | window_end | window_rowtime | item | max_price |
+-------------------------+-------------------------+-------------------------+------+-----------+
| 2021-04-15 14:00:00.000 | 2021-04-15 14:10:00.000 | 2021-04-15 14:09:59.999 | A | 1.8 |
| 2021-04-15 14:00:00.000 | 2021-04-15 14:10:00.000 | 2021-04-15 14:09:59.999 | B | 2.5 |
| 2021-04-15 14:00:00.000 | 2021-04-15 14:10:00.000 | 2021-04-15 14:09:59.999 | C | 3.8 |
+-------------------------+-------------------------+-------------------------+------+-----------+
TIMESTAMP_LTZ의 이벤트 타임 속성
소스의 타임스탬프 데이터가 epoch 시간(보통 long 값, 예: 1618989564564)으로 표현된다면, 이벤트 타임 속성을 TIMESTAMP_LTZ 컬럼으로 정의하는 것을 권장해요.
Flink SQL> CREATE TABLE MyTable3 (
item STRING,
price DOUBLE,
ts BIGINT, -- long time value in epoch milliseconds
ts_ltz AS TO_TIMESTAMP_LTZ(ts, 3),
WATERMARK FOR ts_ltz AS ts_ltz - INTERVAL '10' SECOND
) WITH (
'connector' = 'socket',
'hostname' = '127.0.0.1',
'port' = '9999',
'format' = 'csv'
);
Flink SQL> CREATE VIEW MyView5 AS
SELECT
TUMBLE_START(ts_ltz, INTERVAL '10' MINUTES) AS window_start,
TUMBLE_END(ts_ltz, INTERVAL '10' MINUTES) AS window_end,
TUMBLE_ROWTIME(ts_ltz, INTERVAL '10' MINUTES) as window_rowtime,
item,
MAX(price) as max_price
FROM MyTable3
GROUP BY TUMBLE(ts_ltz, INTERVAL '10' MINUTES), item;
Flink SQL> DESC MyView5;
+----------------+----------------------------+-------+-----+--------+-----------+
| name | type | null | key | extras | watermark |
+----------------+----------------------------+-------+-----+--------+-----------+
| window_start | TIMESTAMP(3) | false | | | |
| window_end | TIMESTAMP(3) | false | | | |
| window_rowtime | TIMESTAMP_LTZ(3) *ROWTIME* | true | | | |
| item | STRING | true | | | |
| max_price | DOUBLE | true | | | |
+----------------+----------------------------+-------+-----+--------+-----------+
MyTable3의 입력 데이터는 다음과 같아요:
A,1.1,1618495260000 # The corresponding utc timestamp is 2021-04-15 14:01:00
B,1.2,1618495320000 # The corresponding utc timestamp is 2021-04-15 14:02:00
A,1.8,1618495380000 # The corresponding utc timestamp is 2021-04-15 14:03:00
B,2.5,1618495440000 # The corresponding utc timestamp is 2021-04-15 14:04:00
C,3.8,1618495500000 # The corresponding utc timestamp is 2021-04-15 14:05:00
C,3.8,1618495860000 # The corresponding utc timestamp is 2021-04-15 14:11:00
Flink SQL> SET 'table.local-time-zone' = 'UTC';
Flink SQL> SELECT * FROM MyView5;
+-------------------------+-------------------------+-------------------------+------+-----------+
| window_start | window_end | window_rowtime | item | max_price |
+-------------------------+-------------------------+-------------------------+------+-----------+
| 2021-04-15 14:00:00.000 | 2021-04-15 14:10:00.000 | 2021-04-15 14:09:59.999 | A | 1.8 |
| 2021-04-15 14:00:00.000 | 2021-04-15 14:10:00.000 | 2021-04-15 14:09:59.999 | B | 2.5 |
| 2021-04-15 14:00:00.000 | 2021-04-15 14:10:00.000 | 2021-04-15 14:09:59.999 | C | 3.8 |
+-------------------------+-------------------------+-------------------------+------+-----------+
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT * FROM MyView5;
UTC 시간대에서의 계산과 비교해 다른 window start, window end, window rowtime을 반환해요.
+-------------------------+-------------------------+-------------------------+------+-----------+
| window_start | window_end | window_rowtime | item | max_price |
+-------------------------+-------------------------+-------------------------+------+-----------+
| 2021-04-15 22:00:00.000 | 2021-04-15 22:10:00.000 | 2021-04-15 22:09:59.999 | A | 1.8 |
| 2021-04-15 22:00:00.000 | 2021-04-15 22:10:00.000 | 2021-04-15 22:09:59.999 | B | 2.5 |
| 2021-04-15 22:00:00.000 | 2021-04-15 22:10:00.000 | 2021-04-15 22:09:59.999 | C | 3.8 |
+-------------------------+-------------------------+-------------------------+------+-----------+
일광 절약 시간 지원 (Daylight Saving Time Support)
Flink SQL은 TIMESTAMP_LTZ 컬럼에 시간 속성을 정의하는 것을 지원하며, 이를 기반으로 Flink SQL은 윈도우 처리에서 TIMESTAMP와 TIMESTAMP_LTZ 타입을 우아하게 사용해 일광 절약 시간(Daylight Saving Time)을 지원해요.
Flink는 타임스탬프 리터럴을 사용해 윈도우를 분할하고, 각 행의 epoch 시간에 따라 데이터에 윈도우를 할당해요. 즉 Flink는 window start와 window end(예: TUMBLE_START와 TUMBLE_END)에 TIMESTAMP 타입을 사용하고, 윈도우 시간 속성(예: TUMBLE_PROCTIME, TUMBLE_ROWTIME)에 TIMESTAMP_LTZ를 사용해요.
텀블 윈도우의 예를 들면, 로스앤젤레스의 일광 절약 시간은 2021-03-14 02:00:00에 시작해요:
long epoch1 = 1615708800000L; // 2021-03-14 00:00:00
long epoch2 = 1615712400000L; // 2021-03-14 01:00:00
long epoch3 = 1615716000000L; // 2021-03-14 03:00:00, skip one hour (2021-03-14 02:00:00)
long epoch4 = 1615719600000L; // 2021-03-14 04:00:00
텀블 윈도우 [2021-03-14 00:00:00, 2021-03-14 00:04:00]는 로스앤젤레스 시간대에서 3시간의 데이터를 수집하지만, DST가 아닌 다른 시간대에서는 4시간의 데이터를 수집해요. 사용자가 해야 할 일은 TIMESTAMP_LTZ 컬럼에 시간 속성을 정의하는 것뿐이에요.
Hop window, Session window, Cumulative window 같은 Flink의 모든 윈도우가 이 방식을 따르며, Flink SQL의 모든 연산이 TIMESTAMP_LTZ를 잘 지원하므로 Flink는 일광 절약 시간대를 우아하게 지원해요.
배치 모드와 스트리밍 모드의 차이 (Difference between Batch and Streaming Mode)
다음 시간 함수들:
- LOCALTIME
- LOCALTIMESTAMP
- CURRENT_DATE
- CURRENT_TIME
- CURRENT_TIMESTAMP
- NOW()
Flink는 이 값들을 실행 모드에 따라 평가해요. 스트리밍 모드에서는 각 레코드마다 평가돼요. 하지만 배치 모드에서는 쿼리가 시작될 때 한 번 평가되고 모든 행에 대해 동일한 결과를 사용해요.
다음 시간 함수들은 배치든 스트리밍이든 상관없이 각 레코드마다 평가돼요:
- CURRENT_ROW_TIMESTAMP()
- PROCTIME()