시간대
시간대 (Time Zone)
Flink는 DATE, TIME, TIMESTAMP, TIMESTAMP_LTZ, INTERVAL YEAR TO MONTH, INTERVAL DAY TO SECOND를 포함한 풍부한 Date/Time 데이터 타입을 제공합니다(자세한 내용은 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는 연, 월, 일, 시, 분, 초 및 소수 초를 나타내는 타임스탬프를 설명합니다.TIMESTAMP는 문자열 리터럴로 지정할 수 있습니다. 예:
Flink SQL> SELECT TIMESTAMP '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)상의 절대 시점을 설명하고, 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;
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT * FROM T1;
TIMESTAMP_LTZ는 시간대 간 비즈니스에 사용할 수 있습니다. 절대 시점(예: 위4001밀리초)이 서로 다른 시간대에서 같은 순간을 설명하기 때문입니다. 같은 시점에 세계의 모든 머신의System.currentTimeMillis()가 같은 값을 반환한다는 배경(예: 위 예제의4001밀리초)이 주어지면, 이것이 절대 시점의 의미입니다.
시간대 사용법 (Time Zone Usage)
로컬 시간대는 현재 세션의 시간대 id를 정의합니다. SQL Client나 애플리케이션에서 시간대를 설정할 수 있습니다.
SQL Client:
-- UTC 시간대로 설정
Flink SQL> SET 'table.local-time-zone' = 'UTC';
-- 상하이 시간대로 설정
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
-- 로스앤젤레스 시간대로 설정
Flink SQL> SET 'table.local-time-zone' = 'America/Los_Angeles';
Java:
EnvironmentSettings envSetting = EnvironmentSettings.inStreamingMode();
TableEnvironment tEnv = TableEnvironment.create(envSetting);
// UTC 시간대로 설정
tEnv.getConfig().setLocalTimeZone(ZoneId.of("UTC"));
// 상하이 시간대로 설정
tEnv.getConfig().setLocalTimeZone(ZoneId.of("Asia/Shanghai"));
// 로스앤젤레스 시간대로 설정
tEnv.getConfig().setLocalTimeZone(ZoneId.of("America/Los_Angeles"));
Scala:
val envSetting = EnvironmentSettings.inStreamingMode()
val tEnv = TableEnvironment.create(envSetting)
// UTC 시간대로 설정
tEnv.getConfig.setLocalTimeZone(ZoneId.of("UTC"))
// 상하이 시간대로 설정
tEnv.getConfig.setLocalTimeZone(ZoneId.of("Asia/Shanghai"))
// 로스앤젤레스 시간대로 설정
tEnv.getConfig.setLocalTimeZone(ZoneId.of("America/Los_Angeles"))
Python:
env_setting = EnvironmentSettings.in_streaming_mode()
t_env = TableEnvironment.create(env_setting)
# UTC 시간대로 설정
t_env.get_config().set_local_timezone("UTC")
# 상하이 시간대로 설정
t_env.get_config().set_local_timezone("Asia/Shanghai")
# 로스앤젤레스 시간대로 설정
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;
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;
Flink SQL> SET 'table.local-time-zone' = 'UTC';
Flink SQL> SELECT * FROM MyView2;
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT * FROM MyView2;
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;
Flink SQL> SELECT * FROM MyView3;
시간 속성과 시간대 (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();
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT PROCTIME();
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;
터미널에서 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;
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT * FROM MyView3;
UTC 시간대에서의 계산과 비교해 다른 window start, window end, window proctime을 반환합니다.
참고: 처리 시간 윈도우는 비결정적이므로 실행할 때마다 다른 윈도우와 다른 집계를 얻습니다. 위 예제는 시간대가 처리 시간 윈도우에 어떤 영향을 주는지 설명하기 위한 것입니다.
이벤트 시간과 시간대 (Event Time and Time Zone)
Flink는 TIMESTAMP 컬럼과 TIMESTAMP_LTZ 컬럼에 이벤트 시간 속성을 정의하는 것을 지원합니다.
TIMESTAMP의 이벤트 시간 속성 (Event Time Attribute on TIMESTAMP)
소스의 타임스탬프 데이터가 연-월-일-시-분-초로 표현된다면(보통 시간대 정보가 없는 문자열 값, 예: 2020-04-15 20:13:40.564), 이벤트 시간 속성을 TIMESTAMP 컬럼으로 정의하는 것이 좋습니다:
Flink SQL> CREATE TABLE MyTable2 (
item STRING,
price DOUBLE,
ts TIMESTAMP(3), -- TIMESTAMP 데이터 타입
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;
터미널에서 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;
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT * FROM MyView4;
UTC 시간대에서의 계산과 비교해 같은 window start, window end, window rowtime을 반환합니다.
TIMESTAMP_LTZ의 이벤트 시간 속성 (Event Time Attribute on TIMESTAMP_LTZ)
소스의 타임스탬프 데이터가 epoch 시간으로 표현된다면(보통 long 값, 예: 1618989564564), 이벤트 시간 속성을 TIMESTAMP_LTZ 컬럼으로 정의하는 것이 좋습니다.
Flink SQL> CREATE TABLE MyTable3 (
item STRING,
price DOUBLE,
ts BIGINT, -- epoch 밀리초 단위의 long 시간 값
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;
MyTable3의 입력 데이터는 다음과 같습니다:
A,1.1,1618495260000 # 해당 UTC 타임스탬프는 2021-04-15 14:01:00
B,1.2,1618495320000 # 해당 UTC 타임스탬프는 2021-04-15 14:02:00
A,1.8,1618495380000 # 해당 UTC 타임스탬프는 2021-04-15 14:03:00
B,2.5,1618495440000 # 해당 UTC 타임스탬프는 2021-04-15 14:04:00
C,3.8,1618495500000 # 해당 UTC 타임스탬프는 2021-04-15 14:05:00
C,3.8,1618495860000 # 해당 UTC 타임스탬프는 2021-04-15 14:11:00
Flink SQL> SET 'table.local-time-zone' = 'UTC';
Flink SQL> SELECT * FROM MyView5;
Flink SQL> SET 'table.local-time-zone' = 'Asia/Shanghai';
Flink SQL> SELECT * FROM MyView5;
UTC 시간대에서의 계산과 비교해 다른 window start, window end, window rowtime을 반환합니다.
서머타임 지원 (Daylight Saving Time Support)
Flink SQL은 TIMESTAMP_LTZ 컬럼에 시간 속성을 정의하는 것을 지원하며, 이를 기반으로 Flink SQL은 윈도우 처리에서 TIMESTAMP와 TIMESTAMP_LTZ 타입을 자연스럽게 사용해 서머타임을 지원합니다.
Flink는 타임스탬프 리터럴로 윈도우를 분할하고 각 행의 epoch 시간에 따라 데이터를 윈도우에 할당합니다. 즉 Flink는 window start와 window end에 TIMESTAMP 타입(예: TUMBLE_START, TUMBLE_END)을 사용하고, 윈도우 시간 속성에는 TIMESTAMP_LTZ(예: TUMBLE_PROCTIME, TUMBLE_ROWTIME)를 사용합니다.
tumble window의 예를 들면, Los_Angeles의 서머타임은 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, 한 시간 건너뜀 (2021-03-14 02:00:00)
long epoch4 = 1615719600000L; // 2021-03-14 04:00:00
tumble window [2021-03-14 00:00:00, 2021-03-14 00:04:00]는 Los_Angeles 시간대에서 3시간의 데이터를 수집하지만, 다른 비-DST 시간대에서는 4시간의 데이터를 수집합니다. 사용자가 해야 할 일은 TIMESTAMP_LTZ 컬럼에 시간 속성을 정의하는 것뿐입니다.
Flink의 모든 윈도우(Hop window, Session window, Cumulative window)는 이 방식을 따르며, 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()