Flink SQL mendukung fungsi jendela TUMBLE, HOP, dan SESSION berdasarkan event time atau processing time.
Fungsi jendela
Flink SQL mendukung agregasi atas jendela tak hingga tanpa perlu secara eksplisit mendefinisikan jendela dalam pernyataan SQL. Flink SQL juga mendukung agregasi atas jendela tertentu. Misalnya, untuk menghitung jumlah pengguna yang mengklik URL dalam satu menit terakhir, definisikan jendela yang mengumpulkan data klik selama satu menit terakhir dan menghitung hasilnya.
Flink SQL mendukung window aggregates dan over aggregates. Topik ini membahas window aggregates, yang menggunakan dua atribut waktu—event time dan processing time—serta mendukung fungsi jendela TUMBLE, HOP, dan SESSION.
Agregasi jendela (TUMBLE, HOP, dan SESSION) yang dikombinasikan dengan fungsi LAST_VALUE, FIRST_VALUE, atau TopN dapat menghasilkan output yang tidak akurat karena mekanisme pemicu jendela atau latensi.
Atribut waktu
Flink SQL mendukung dua atribut waktu: event time dan processing time. Perilaku jendela berbeda tergantung pada atribut yang digunakan.
-
Event time: biasanya merupakan timestamp yang tertanam dalam suatu catatan.
Jendela ditutup ketika watermark melebihi waktu akhir jendela. Output hanya dihasilkan ketika data penutup tiba. Untuk subtask tunggal, watermark meningkat secara monotonik. Dengan beberapa subtask atau tabel sumber, Flink menggunakan nilai watermark minimum.
Penting-
Jika terdapat catatan yang tidak berurutan atau subtask/partisi tidak memiliki data, watermark tidak dapat maju dan jendela mungkin tidak tertutup. Untuk mengatasi hal ini, tentukan offset watermark untuk data yang tidak berurutan, dan pastikan semua subtask serta partisi menerima aliran data. Jika suatu partisi menganggur, tambahkan
table.exec.source.idle-timeout: 10ske bidang Other Configuration pada bagian Parameters di tab Configuration halaman Penyebaran. Lihat Configuration untuk detail parameter. -
Setelah data diproses menggunakan GROUP BY, operasi JOIN antara dua aliran data, atau node OVER window, properti watermark hilang dan event time tidak dapat lagi digunakan untuk pembuatan jendela.
-
-
Processing time: waktu jam sistem saat Flink memproses suatu event.
Processing time dihasilkan oleh Flink dan tidak ada dalam data mentah Anda, sehingga Anda harus secara eksplisit mendefinisikan kolom processing time.
CatatanKarena processing time bergantung pada kecepatan kedatangan event dan urutan pemrosesan, hasil backtrack dapat berbeda antar eksekusi.
Agregasi jendela bertingkat
Setelah agregasi jendela selesai, kolom rowtime kehilangan atribut event time-nya. Gunakan fungsi pembantu seperti TUMBLE_ROWTIME, HOP_ROWTIME, atau SESSION_ROWTIME untuk mengambil max(rowtime) dari jendela dan menggunakannya sebagai rowtime baru. Nilai yang dikembalikan sama dengan window_end - 1, bertipe TIMESTAMP, dan mempertahankan atribut rowtime. Sebagai contoh, untuk jendela [00:00, 00:15), nilai 00:14:59.999 dikembalikan.
Contoh berikut melakukan agregasi jendela bergulir (tumbling) 1 jam di atas agregasi jendela bergulir 1 menit:
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 -- Definisikan watermark untuk 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 hanya mengizinkan ekspor pernyataan DDL bertipe VARCHAR. Oleh karena itu, DataHub digunakan untuk menyimpan data.
...
);
CREATE TEMPORARY VIEW one_minute_window_output AS
SELECT
TUMBLE_ROWTIME(ts, INTERVAL '1' MINUTE) as rowtime, -- Gunakan TUMBLE_ROWTIME sebagai waktu agregasi untuk jendela tingkat kedua.
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;
Hasil antara
Data antara jendela terdiri dari keyed state dan data timer, yang masing-masing disimpan di backend berbeda. Pilih kombinasi berdasarkan karakteristik pekerjaan Anda:
|
Penyimpanan untuk state berkunci |
Penyimpanan timer |
|
Memory |
|
|
Memory |
|
|
Memory |
|
|
File |
Timer terutama digunakan untuk memicu jendela yang telah kedaluwarsa. Simpan timer di memory untuk performa terbaik. Jika jumlah timer sangat banyak atau kapasitas memory terbatas, gunakan RocksDBStateBackend untuk menyimpannya dalam file RocksDB.