Flink Time Characteristic und AutoWatermarkInterval
setAutoWatermarkInterval(interval)Erzeugt in Apache Flink Wasserzeichen für nachgeschaltete Bediener, damit diese ihre Ereigniszeit verlängern können.
Wenn das Wasserzeichen während des angegebenen Intervalls nicht geändert wurde (keine Ereignisse eingetroffen), gibt die Laufzeit keine Wasserzeichen aus? Wenn andererseits ein neues Ereignis vor dem nächsten Intervall eintrifft, wird sofort ein neues Wasserzeichen ausgegeben oder es wird in die Warteschlange gestellt / gewartet, bis das nächste Intervall für setAutoWatermarkInterval erreicht ist.
Ich bin gespannt auf die beste Konfiguration von AutoWatermarkInterval (insbesondere für Quellen mit hoher Rate): Je kleiner dieser Wert ist, desto geringer ist die Verzögerung zwischen Verarbeitungszeit und Ereigniszeit, jedoch aufgrund des höheren BW-Verbrauchs beim Senden der Wasserzeichen . Ist das wahr richtig?
Wenn ich dagegen env.setStreamTimeCharacteristic (TimeCharacteristic.IngestionTime) verwendet habe, weist die Flink-Laufzeit automatisch Zeitstempel und Wasserzeichen zu (Zeitstempel entsprechen dem Zeitpunkt, zu dem das Ereignis in die Flink-Datenfluss-Pipeline eingegeben wurde, dh den Quelloperator), obwohl dies auch bei ingestionTime möglich ist Definieren Sie weiterhin einen Verarbeitungszeit-Timer (in der processElement-Funktion) wie unten gezeigt:
long timer = context.timestamp() + Timeout.
context.timerService().registerProcessingTimeTimer(timer);
Dabei ist context.timestamp () die von Flink festgelegte Aufnahmezeit.
Vielen Dank.
Antworten
Das autoWatermarkInterval betrifft nur Wasserzeichengeneratoren, die darauf achten. Sie haben auch die Möglichkeit, ein Wasserzeichen in Kombination mit der Ereignisverarbeitung zu generieren.
Für diejenigen Wasserzeichengeneratoren, die das autoWatermarkInterval verwenden (was definitiv der Normalfall ist), sammeln sie Beweise dafür, was das nächste Wasserzeichen sein sollte, als Nebeneffekt beim Zuweisen von Zeitstempeln für jedes Ereignis. Wenn ein Timer ausgelöst wird (basierend auf dem autoWatermarkInterval), wird der Wasserzeichengenerator von der Flink-Laufzeit aufgefordert, das nächste Wasserzeichen zu erstellen. Das Wasserzeichen wartete nicht irgendwo und wurde auch nicht in die Warteschlange gestellt, sondern wird bei Bedarf auf der Grundlage von Informationen erstellt, die vom Zeitstempelzuweiser gespeichert wurden. Dies ist normalerweise der maximale Zeitstempel, der bisher im Stream angezeigt wurde.
Ja, häufigere Wasserzeichen bedeuten mehr Aufwand für die Kommunikation und Verarbeitung sowie eine geringere Latenz. Sie müssen anhand der Anforderungen Ihrer Anwendung entscheiden, wie dieser Kompromiss zwischen Durchsatz und Latenz behandelt werden soll.
Sie können unabhängig von der TimeCharacteristic immer Zeitgeber für die Verarbeitungszeit verwenden. (Übrigens lösen Wasserzeichen auf niedriger Ebene nur Ereigniszeit-Timer aus, sei es in Prozessfunktionen, Fenstern usw.)