Flink Time Characteristic dan AutoWatermarkInterval

Aug 29 2020

Di Apache Flink, setAutoWatermarkInterval(interval)menghasilkan tanda air ke operator hilir sehingga mereka memajukan waktu acara mereka.

Jika tanda air belum diubah selama interval yang ditentukan (tidak ada peristiwa yang tiba), runtime tidak akan mengeluarkan tanda air? Di sisi lain, jika acara baru tiba sebelum interval berikutnya, tanda air baru akan segera dikeluarkan atau akan diantrekan / menunggu hingga interval setAutoWatermarkInterval berikutnya tercapai.

Saya ingin tahu tentang konfigurasi apa yang terbaik AutoWatermarkInterval (terutama untuk sumber tingkat tinggi): semakin banyak nilai ini kecil, semakin banyak jeda antara waktu pemrosesan dan waktu acara akan kecil, tetapi pada overhead penggunaan BW lebih banyak untuk mengirim tanda air . Apakah itu benar akurat?

Di sisi lain, Jika saya menggunakan env.setStreamTimeCharacteristic (TimeCharacteristic.IngestionTime), runtime Flink akan secara otomatis menetapkan stempel waktu dan tanda air (stempel waktu sesuai dengan waktu saat peristiwa memasuki pipa aliran data Flink yaitu operator sumber), meskipun demikian bahkan dengan ingestionTime kita bisa masih menentukan timer waktu pemrosesan (dalam fungsi processElement) seperti yang ditunjukkan di bawah ini:

long timer = context.timestamp() + Timeout.
context.timerService().registerProcessingTimeTimer(timer);

di mana context.timestamp () adalah waktu penyerapan yang disetel oleh Flink.

Terima kasih.

Jawaban

2 DavidAnderson Aug 29 2020 at 17:13

AutoWatermarkInterval hanya memengaruhi generator watermark yang memperhatikannya. Mereka juga memiliki kesempatan untuk menghasilkan tanda air yang dikombinasikan dengan pemrosesan acara.

Untuk generator watermark yang menggunakan autoWatermarkInterval (yang tentunya merupakan kasus normal), mereka mengumpulkan bukti tentang apa yang seharusnya menjadi watermark berikutnya sebagai efek samping dari menetapkan cap waktu untuk setiap kejadian. Ketika pengatur waktu menyala (berdasarkan autoWatermarkInterval), generator tanda air kemudian diminta oleh runtime Flink untuk menghasilkan tanda air berikutnya. Tanda air tidak menunggu di suatu tempat, juga tidak diantrekan, melainkan dibuat sesuai permintaan, berdasarkan informasi yang telah disimpan oleh pemberi stempel waktu - yang biasanya merupakan stempel waktu maksimum yang terlihat sejauh ini di aliran.

Ya, watermark yang lebih sering berarti lebih banyak overhead untuk berkomunikasi dan memprosesnya, dan latensi yang lebih rendah. Anda harus memutuskan cara menangani tradeoff throughput / latensi ini berdasarkan persyaratan aplikasi Anda.

Anda selalu dapat menggunakan timer waktu pemrosesan, terlepas dari TimeCharacteristic tersebut. (Ngomong-ngomong, pada level rendah, satu-satunya hal yang dilakukan watermark adalah memicu pengatur waktu acara, baik itu dalam fungsi proses, jendela, dll.)