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

  1. Registration. The computation calls output.add_timer / output.addTimer / output->AddTimer, passing TriggerTimestamp and EventTimestamp.
  2. Storage. The timer is reliably stored in YTsaurus and additionally cached in the process memory.
  3. Triggering. When the EventWatermark (or another configured watermark) exceeds TriggerTimestamp, the timer is delivered to the computation.
  4. 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

time_type

Type: NYT::NFlow::ETimeType
Default value: event_time
The time type used by the timer: event_time, system_time, real_time

streams

Type: std::optional<THashSet<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>>>
List of streams that the timer monitors.

streams_with_delays

Type: std::optional<THashMap<NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>, TDuration>>
List of streams that the timer monitors with an individual delay specified

deduplicate_equal_timestamps

Type: bool
Default value: true
Enables timer deduplication: if you try to create a timer with the same key and trigger_timestamp, only one timer remains — the one with the smaller EventTimestamp value

Explanations:

  • time_type defines the scale that the system uses to compare TriggerTimestamp with the current watermark:

    • event_time (default) — compared with EventWatermark across all input streams (or the streams listed in streams).
    • system_time — compared with SystemWatermark.
    • 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_delays also lets you set an individual delay for each stream.

  • deduplicate_equal_timestamps — merges timers in the same stream with the same key and TriggerTimestamp, keeping the smallest EventTimestamp. 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

message_id

Type: NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TMessageIdTag>
Required parameter
Unique timer ID.

system_timestamp

Type: NYT::TStrongTypedef<unsigned long, NYT::NFlow::TSystemTimestampTag, NYT::TStrongTypedefOptions{true}>
Required parameter
Timestamp of the timer creation.

event_timestamp

Type: NYT::TStrongTypedef<unsigned long, NYT::NFlow::TSystemTimestampTag, NYT::TStrongTypedefOptions{true}>
Required parameter
Timestamp of the real event associated with this timer.

stream_id

Type: NYT::NFlow::TStrongIdentifierTypedef<NYT::NFlow::TStreamIdTag>
Required parameter
The stream this message belongs to.

key

Type: NYT::TStrongTypedef<NYT::NFlow::TCompactUnversionedOwningRow, NYT::NFlow::TKeyTag, NYT::TStrongTypedefOptions{true}>
Required parameter
Key of the timer.

key_schema

Type: NYT::TIntrusivePtr<NYT::NTableClient::TTableSchema>
Required parameter
Key schema, matches the group_by_schema of the corresponding Computation.

trigger_timestamp

Type: NYT::TStrongTypedef<unsigned long, NYT::NFlow::TSystemTimestampTag, NYT::TStrongTypedefOptions{true}>
Required parameter
Time for the timer to fire.

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).

See also