///|
pub(all) struct WatermarkState {
policy : WatermarkPolicy
max_event_time_ms : Int
current_watermark_ms : Int
last_arrival_ms : Int
initialized : Bool
} derive(Eq, Debug)
///|
pub fn WatermarkState::new(policy : WatermarkPolicy) -> WatermarkState {
{
policy,
max_event_time_ms: 0,
current_watermark_ms: 0,
last_arrival_ms: 0,
initialized: false,
}
}
///|
pub fn WatermarkState::observe(
self : WatermarkState,
event : StreamEvent,
) -> WatermarkState {
let max_event_time = if !self.initialized ||
event.event_time_ms > self.max_event_time_ms {
event.event_time_ms
} else {
self.max_event_time_ms
}
let next_watermark = max_event_time - self.policy.max_out_of_order_ms
{
policy: self.policy,
max_event_time_ms: max_event_time,
current_watermark_ms: if self.initialized &&
self.current_watermark_ms > next_watermark {
self.current_watermark_ms
} else {
next_watermark
},
last_arrival_ms: event.arrival_time_ms,
initialized: true,
}
}
///|
pub fn WatermarkState::mark_idle(self : WatermarkState, now_ms : Int) -> Bool {
self.initialized &&
now_ms >= self.last_arrival_ms + self.policy.idle_timeout_ms
}
///|
pub fn WatermarkState::classify(
self : WatermarkState,
event : StreamEvent,
) -> EventStatus {
if !self.initialized || event.event_time_ms >= self.current_watermark_ms {
OnTime
} else if event.event_time_ms + self.policy.allowed_lateness_ms >=
self.current_watermark_ms {
Late
} else {
TooLate
}
}
///|
pub fn WatermarkState::deadline_for(
self : WatermarkState,
event_time_ms : Int,
) -> Int {
event_time_ms + self.policy.allowed_lateness_ms
}