---
metadata:
  - name: generator
    content: Diplodoc Platform v5.50.6
alternate:
  - https://ytsaurus.tech/docs/en/flow/concepts/timers.md
  - https://ytsaurus.tech/docs/ru/flow/concepts/timers.md
---
> **Documentation Index:** Fetch the complete configuration index at https://ytsaurus.tech/docs/en/llms.txt

<!-- source: en/_includes/flow/concepts/timers.md -->
# Timers in YTsaurus Flow

## Why you need timers {#why-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](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#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 {#how-timers-work}

### Lifecycle {#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](https://ytsaurus.tech/docs/en/flow/concepts/watermarks.md#event-watermark) (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](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#deduplication)).

### Timer fields {#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](#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](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#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 {#supported-computations}

| Process-function adapter | Timer support |
|---|---|
| `TProcessFunctionComputation` | ✓ |
| `TProcessFunctionSwiftMapComputation` | ✗ |
| `TProcessFunctionSourceComputation` | ✗ |
| `TProcessFunctionTransformOrderedSourceComputation` | ✗ |

## Configuration {#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: en/flow/generated_docs/NYT_NFlow_TTimerSpec.md -->
<!-- This file is generated by yt/yt/flow/yandex/tools/generate_yson_struct_doc/generate.sh script -->
<!-- Before using doc generation tool check readme: yt/yt/flow/yandex/tools/generate_yson_struct_doc/README.md -->
Source: [yt/yt/flow/library/cpp/common/spec.h](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/cpp/common/spec.h)

#|
|| **Parameter** | **Description** ||
|| `time_type` | **Type**: [NYT::NFlow::ETimeType](https://ytsaurus.tech/docs/en/flow/generated_docs/all_yson_structs.md#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](https://ytsaurus.tech/docs/en/flow/generated_docs/all_yson_structs.md#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 ||
|#
<!-- endsource: en/flow/generated_docs/NYT_NFlow_TTimerSpec.md -->

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 {#timer-ordering}

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 {#timer-structure}

<!-- source: en/flow/generated_docs/NYT_NFlow_TTimerSerializer.md -->
<!-- This file is generated by yt/yt/flow/yandex/tools/generate_yson_struct_doc/generate.sh script -->
<!-- Before using doc generation tool check readme: yt/yt/flow/yandex/tools/generate_yson_struct_doc/README.md -->
Source: [yt/yt/flow/library/cpp/common/timer-inl.h](https://github.com/ytsaurus/ytsaurus/tree/main/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. ||
|#
<!-- endsource: en/flow/generated_docs/NYT_NFlow_TTimerSerializer.md -->

## API by language {#api}

### C++ {#api-cpp}

```cpp
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++)](https://ytsaurus.tech/docs/en/flow/cpp/process-functions.md).

### Java {#api-java}

```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)](https://ytsaurus.tech/docs/en/flow/java/computation.md) section.

### Python {#api-python}

```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)](https://ytsaurus.tech/docs/en/flow/python/computation.md) section.

### Go {#api-go}

```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)](https://ytsaurus.tech/docs/en/flow/go/computation.md) section.

## Examples {#examples}

Example implementation of a join with wait (impression + click, 30-minute timeout):

- [C++](https://ytsaurus.tech/docs/en/flow/cpp/examples/wait_click_join.md)
- [Java](https://ytsaurus.tech/docs/en/flow/java/examples/wait_click_join.md)
- [Python](https://ytsaurus.tech/docs/en/flow/python/examples/wait_click_join.md)
- [Go](https://ytsaurus.tech/docs/en/flow/go/examples/wait_click_join.md)

## Limitations and known issues {#limitations}

{% note 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.

{% endnote %}

- 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

- [Watermarks and Timers](https://ytsaurus.tech/docs/en/flow/concepts/watermarks.md)
- [Stateful processing](https://ytsaurus.tech/docs/en/flow/concepts/stateful.md)
- [Computation (C++)](https://ytsaurus.tech/docs/en/flow/cpp/computation.md)
- [Computation (Java)](https://ytsaurus.tech/docs/en/flow/java/computation.md)
- [Computation (Python)](https://ytsaurus.tech/docs/en/flow/python/computation.md)
- [Computation (Go)](https://ytsaurus.tech/docs/en/flow/go/computation.md)
<!-- endsource: en/_includes/flow/concepts/timers.md -->