Flink SQL は、イベント時間または処理時間に基づく TUMBLE、HOP、SESSION のウィンドウ関数をサポートしています。
ウィンドウ関数
Flink SQL は、SQL 文でウィンドウを明示的に定義せずに無限ウィンドウに対する集約をサポートしています。また、特定のウィンドウに対する集約もサポートしています。たとえば、過去 1 分間に URL をクリックしたユーザーをカウントするには、過去 1 分間のクリックデータを収集し、結果を計算するためのウィンドウを定義します。
Flink SQL は、ウィンドウ集計と OVER 集計をサポートしています。本トピックでは、ウィンドウ集計について説明します。ウィンドウ集計は、イベント時間と処理時間という 2 つの時間属性を使用し、TUMBLE、HOP、SESSION のウィンドウ関数をサポートしています。
ウィンドウ集計 (TUMBLE、HOP、SESSION) を LAST_VALUE、FIRST_VALUE、または TopN 関数と組み合わせて使用すると、ウィンドウのトリガーメカニズムやレイテンシーにより、不正確な結果が生じる可能性があります。
時間属性
Flink SQL は、イベント時間と処理時間という 2 つの時間属性をサポートしています。使用する属性によって、ウィンドウの動作が異なります。
-
イベント時間:通常、レコードに埋め込まれたタイムスタンプです。
ウォーターマークがウィンドウの終了時刻を超えると、ウィンドウが閉じます。出力は、ウィンドウが閉じたときにのみ生成されます。
重要-
順序が乱れたレコードが存在する場合、またはサブタスクやパーティションにデータがない場合、ウォーターマークが進まず、ウィンドウが閉じない可能性があります。この問題を解決するには、順序が乱れたデータに対してウォーターマークオフセットを指定し、すべてのサブタスクとパーティションにデータが流れるようにしてください。パーティションがアイドル状態の場合は、Deployments ページの **[Configuration]** タブの **[Parameters]** にある **[Other Configuration]** フィールドに
table.exec.source.idle-timeout: 10sを追加します。パラメーターの詳細については、「Configuration」をご参照ください。 -
GROUP BY、2 つのデータストリームの JOIN 操作、または OVER ウィンドウノードを使用してデータが処理された後、ウォーターマークプロパティが失われ、イベント時間をウィンドウ処理に使用できなくなります。
-
-
処理時間:Flink がイベントを処理するときのシステムクロック時刻です。
処理時間は Flink によって生成され、生データには存在しないため、処理時間列を明示的に定義する必要があります。
説明処理時間はイベントの到着速度と処理順序に依存するため、バックトラック結果は実行ごとに異なる可能性があります。
カスケードウィンドウ集計
ウィンドウ集計が完了すると、rowtime 列はそのイベント時間属性を失います。ヘルパー関数 (たとえば TUMBLE_ROWTIME、HOP_ROWTIME、または SESSION_ROWTIME) を使用して、ウィンドウの max(rowtime) を取得し、それを新しい rowtime として使用します。返される値は window_end - 1 と等しく、TIMESTAMP 型であり、rowtime 属性を保持します。たとえば、ウィンドウ [00:00, 00:15) の場合、値 00:14:59.999 が返されます。
次の例は、1 分間のタンブリングウィンドウに対して 1 時間のタンブリングウィンドウをカスケードする方法を示しています。
CREATE TEMPORARY TABLE user_clicks(
username varchar,
click_url varchar,
eventtime varchar,
ts AS TO_TIMESTAMP(eventtime),
WATERMARK FOR ts AS ts - INTERVAL '2' SECOND -- rowtime のウォーターマークを定義します。
) with (
'connector'='sls',
...
);
CREATE TEMPORARY TABLE tumble_output(
window_start TIMESTAMP,
window_end TIMESTAMP,
username VARCHAR,
clicks BIGINT
) with (
'connector'='datahub' -- Simple Log Service は VARCHAR 型のデータしかエクスポートできないため、DataHub を使用してデータを保存します。
...
);
CREATE TEMPORARY VIEW one_minute_window_output AS
SELECT
TUMBLE_ROWTIME(ts, INTERVAL '1' MINUTE) as rowtime, -- レベル 2 ウィンドウの集約時間として TUMBLE_ROWTIME を使用します。
username,
COUNT(click_url) as cnt
FROM user_clicks
GROUP BY TUMBLE(ts, INTERVAL '1' MINUTE),username;
BEGIN statement set;
INSERT INTO tumble_output
SELECT
TUMBLE_START(rowtime, INTERVAL '1' HOUR),
TUMBLE_END(rowtime, INTERVAL '1' HOUR),
username,
SUM(cnt)
FROM one_minute_window_output
GROUP BY TUMBLE(rowtime, INTERVAL '1' HOUR), username;
END;
中間結果
ウィンドウの中間データは、キード付き状態とタイマーデータで構成され、それぞれ異なるバックエンドに保存されます。ジョブの特性に基づいて組み合わせを選択してください。
|
キー付き状態のストレージ |
タイマーのストレージ |
|
メモリ |
|
|
メモリ |
|
|
メモリ |
|
|
ファイル |
タイマーは主に期限切れのウィンドウをトリガーするために使用されます。最高のパフォーマンスを得るには、タイマーをメモリに保存してください。タイマーが多数ある場合やメモリが限られている場合は、RocksDBStateBackend を使用して RocksDB ファイルに保存してください。