Timers in YTsaurus Flow
Why you need timers
Many stream processing tasks require waiting. For example, you might need to treat a conversion as failed if no click arrives within 30 minutes after an impression. The standard message exchange between computations isn’t suitable for this: a message is either present or not, and you can’t wait for an event’s “absence”.
Timers solve this problem. A process function in transform mode can create a timer — tell the system, “Wake me up when time X arrives.” When the moment comes, the function receives a ProcessTimer / on_timer / onTimer call and can make a decision based on the accumulated state.
Typical use cases:
- Join with wait: correlate an ad impression with a click that might arrive with a delay.
- Timeout: handle a situation where the expected event never arrives.
- Windowed aggregations: close a time window and output the result when the accumulated data becomes stale.
How timers work
Lifecycle
- Registration. The computation calls
output.add_timer/output.addTimer/output->AddTimer, passingTriggerTimestampandEventTimestamp. - Storage. The timer is reliably stored in YTsaurus and additionally cached in the process memory.
- Triggering. When the EventWatermark (or another configured watermark) exceeds
TriggerTimestamp, the timer is delivered to the computation. - Processing and deletion. The processing result and the timer deletion are recorded in a single transaction — this ensures exactly-once guarantees (see Deduplication).
Timer fields
Each timer includes two timestamp fields:
TriggerTimestamp— the moment when the timer should trigger. The time scale (event_time,system_time,real_time) is set in the configuration.EventTimestamp— the business time of the original event, held in the timer. It lets you know when the triggering event occurred while processing the timer.
The timer is also bound to the grouping key (the same key that Computation uses), so it automatically lands in the correct partition and stays isolated from other keys.
Which computations support timers
| Process-function adapter | Timer support |
|---|---|
TProcessFunctionComputation |
✓ |
TProcessFunctionSwiftMapComputation |
✗ |
TProcessFunctionSourceComputation |
✗ |
TProcessFunctionTransformOrderedSourceComputation |
✗ |
Configuration
To enable timers in a computation, you must fill the timers field in its spec. Each array element describes a single timer stream with the following parameters:
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter |
Description |
|
|
Type: NYT::NFlow::ETimeType |
|
|
Type: |
|
|
Type: |
|
|
Type: |
Explanations:
-
time_typedefines the scale that the system uses to compareTriggerTimestampwith the current watermark:event_time(default) — compared withEventWatermarkacross allinputstreams (or the streams listed instreams).system_time— compared withSystemWatermark.real_time— compared with real astronomical time.
-
streams/streams_with_delays— let you limit which input streams contribute to the watermark calculation for this timer.streams_with_delaysalso lets you set an individual delay for each stream. -
deduplicate_equal_timestamps— merges timers in the same stream with the same key andTriggerTimestamp, keeping the smallestEventTimestamp. Deduplication is not guaranteed: duplicates are possible. Enabled by default.
Processing order
Within one timer stream, timers for the same key are processed in TriggerTimestamp order.
Processing order across timer streams is not guaranteed. To make timer A fire after timer B, add stream B to the streams field in timer A's spec. Timer A will then wait for stream B's watermark to advance.
Timer structure
Source: yt/yt/flow/library/cpp/common/timer-inl.h
|
Parameter |
Description |
|
|
Type: |
|
|
Type: |
|
|
Type: |
|
|
Type: |
|
|
Type: |
|
|
Type: |
|
|
Type: |
API by language
C++
void ProcessMessage(
const TInputMessageConstPtr& message,
const IOutputCollectorPtr& output,
const IRuntimeContextPtr& /*context*/) override
{
output->AddTimer(
TSystemTimestamp(message->EventTimestamp.Underlying() + TDuration::Minutes(30).Seconds()),
message->EventTimestamp);
}
void ProcessTimer(
const TInputTimerConstPtr& timer,
const IOutputCollectorPtr& output,
const IRuntimeContextPtr& context) override
{
// timer->Key, timer->EventTimestamp, timer->TriggerTimestamp
auto builder = context->MakeOutputMessageBuilder();
// ...
output->AddMessage(builder.Finish());
}
For more details, see Process functions (C++).
Java
// Create a timer in onMessage:
output.addTimer(message.getEventTimestamp() + 30 * 60_000_000_000L, message.getEventTimestamp());
// Process the triggered timer (RowFunction):
@Override
public void onTimer(Timer timer, OutputCollector output, RuntimeContext ctx) {
// timer.getKey(), timer.getEventTimestamp(), timer.getTriggerTimestamp()
}
// For BatchFunction:
@Override
public void onTimers(List<Timer> timers, OutputCollector output, RuntimeContext ctx) { ... }
For more details, see the Computation (Java) section.
Python
# Create a timer in on_message:
output.add_timer(trigger_timestamp=message.event_timestamp + 30 * 60_000_000_000, event_timestamp=message.event_timestamp)
# Process the triggered timer (RowFunction):
def on_timer(self, timer, output, ctx):
# timer.key, timer.event_timestamp, timer.trigger_timestamp, timer.stream_id
# For BatchFunction:
def on_timers(self, timers, output, ctx):
for timer in timers:
...
For more details, see the Computation (Python) section.
Go
// Create a timer in OnMessage:
out.AddTimer(flow.TimerRequest{
TriggerTimestamp: msg.EventTimestamp + 30*60_000_000_000,
EventTimestamp: msg.EventTimestamp,
})
// Process the triggered timer (RowFunction):
func (f myFunction) OnTimer(
ctx context.Context,
rt flow.Runtime,
timer flow.Timer,
out flow.OutputCollector,
) error {
// timer.Key, timer.EventTimestamp, timer.TriggerTimestamp, timer.StreamID
return nil
}
// For BatchFunction:
func (f myFunction) OnTimers(ctx context.Context, rt flow.Runtime, timers []flow.Timer, out flow.OutputCollector) error { ... }
For more details, see the Computation (Go) section.
Examples
Example implementation of a join with wait (impression + click, 30-minute timeout):
Limitations and known issues
Warning
Process memory. In the current implementation, all active timers are additionally stored in the process memory, on top of YTsaurus. With a large number of timers, this can lead to Out of Memory errors when jobs start.
- Swift-map mode doesn’t support timers. If you need timer functionality, run the process function under
TProcessFunctionComputation. - Timestamps are transmitted in nanoseconds (uint64).