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

<!-- source: ru/_includes/flow/concepts/ordering.md -->
# Порядок обработки сообщений в YTsaurus Flow

В распределённой потоковой системе строгий глобальный порядок всех сообщений обеспечить непрактично: данные поступают из множества источников, обрабатываются параллельно в разных [партициях](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#partition) и [компьютейшенах](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#computation). Flow решает эту задачу иначе — через механизм `AlignmentTimestamp`.

Каждое сообщение несёт [AlignmentTimestamp](#alignment-timestamp) — монотонно возрастающую временну́ю метку, которая определяет приоритет обработки:

- **Внутри одного [стрима](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#stream)**: сообщения обрабатываются в порядке `(AlignmentTimestamp, message_id)` — полный детерминированный порядок.
- **Между разными [стримами](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#stream)**: используется сливающая очередь приоритетов по [StabilizedEventTimestamp](#stabilized-event-timestamp) (вычисляется на основе `AlignmentTimestamp`) с поправкой на [stream_delays](https://ytsaurus.tech/docs/ru/flow/concepts/spec.md#inputordering); при равенстве — по `TaskId`.

Помимо приоритизации существует гарантия порядка для **производных сообщений** — подробнее в разделе [Гарантии порядка](#ordering-guarantees).

## AlignmentTimestamp {#alignment-timestamp}

Помимо [SystemTimestamp и EventTimestamp](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#timestamps-and-watermarks), каждое сообщение во Flow несёт поле `AlignmentTimestamp` — временную метку, используемую для выравнивания прогресса обработки между потоками и партициями.

Правила выставления `AlignmentTimestamp`:

- **Сообщения из source-потоков** — равен `WriteTimestamp`, то есть времени записи сообщения в персистентную очередь (например, в QYT).
- **Выходные сообщения `TransformComputation`** — равен `SystemTimestamp` на момент создания сообщения.
- **Во всех остальных случаях** — наследуется от родительских сообщений без изменений.

## StabilizedEventTimestamp {#stabilized-event-timestamp}

`StabilizedEventTimestamp` &mdash; это вычислимый таймстемп, равный `AlignmentTimestamp` + `bias`, где `bias` это средняя разность между `EventTimestamp` и `AlignmentTimestamp` в рамках стрима. `bias` считается на основе тех сообщений стрима, что сейчас находятся в обработке. Используется как неубывающая замена для `EventTimestamp` для приоритезации обработки сообщений из разных стримов.

## Гарантии порядка {#ordering-guarantees}

Порядок сообщений в потоке в общем случае не гарантируется, однако для производных сообщений существует следующая гарантия.

Если во входном потоке сообщение A предшествует сообщению B, и у них есть производные сообщения A' и B' такие, что:

- ключи в их [lineage](https://ytsaurus.tech/docs/ru/flow/concepts/lineage.md) совпадают (то есть [ключи группировки](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#key) во всех промежуточных компьютейшенах по пути от источника до текущего компьютейшена одинаковы),
- и в текущем компьютейшене у A' и B' один и тот же ключ,

то A' будет обработано раньше, чем B'.

## Приоритизация партиций внутри одного стрима {#partition-prioritization}

Внутри одного source-стрима партиции приоритизируются по `WriteTimestamp` их сообщений: партиции с более ранними сообщениями обрабатываются в первую очередь. Приоритизация реализована через порядок обработки выходных сообщений в следующих компьютейшенах.

## Приоритизация между стримами {#cross-stream-prioritization}

Между разными стримами приоритизация выполняется по `EventTimestamp`, но опосредованно: фактически упорядочивание происходит по `AlignmentTimestamp` с поправкой, специфичной для каждого стрима.

Поправка вычисляется как среднее отличие между `EventTimestamp` и `AlignmentTimestamp` среди сообщений, находящихся в данный момент во всех выходных буферах по данному стриму:

$$ordering\_priority = AlignmentTimestamp + avg(EventTimestamp − AlignmentTimestamp)$$

Такой подход позволяет учитывать специфику каждого стрима (например, систематическую задержку между временем записи в очередь и временем события) и обеспечивает более справедливое межстримовое упорядочивание.

## См. также

- [Computation](https://ytsaurus.tech/docs/ru/flow/concepts/computation.md)
- [Вотермарки](https://ytsaurus.tech/docs/ru/flow/concepts/watermarks.md)
- [Таймеры](https://ytsaurus.tech/docs/ru/flow/concepts/timers.md)
<!-- endsource: ru/_includes/flow/concepts/ordering.md -->
