Порядок обработки сообщений в YTsaurus Flow

В распределённой потоковой системе строгий глобальный порядок всех сообщений обеспечить непрактично: данные поступают из множества источников, обрабатываются параллельно в разных партициях и компьютейшенах. Flow решает эту задачу иначе — через механизм AlignmentTimestamp.

Каждое сообщение несёт AlignmentTimestamp — монотонно возрастающую временну́ю метку, которая определяет приоритет обработки:

  • Внутри одного стрима: сообщения обрабатываются в порядке (AlignmentTimestamp, message_id) — полный детерминированный порядок.
  • Между разными стримами: используется сливающая очередь приоритетов по StabilizedEventTimestamp (вычисляется на основе AlignmentTimestamp) с поправкой на stream_delays; при равенстве — по TaskId.

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

AlignmentTimestamp

Помимо SystemTimestamp и EventTimestamp, каждое сообщение во Flow несёт поле AlignmentTimestamp — временную метку, используемую для выравнивания прогресса обработки между потоками и партициями.

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

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

StabilizedEventTimestamp

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

Гарантии порядка

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

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

  • ключи в их lineage совпадают (то есть ключи группировки во всех промежуточных компьютейшенах по пути от источника до текущего компьютейшена одинаковы),
  • и в текущем компьютейшене у A' и B' один и тот же ключ,

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

Приоритизация партиций внутри одного стрима

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

Приоритизация между стримами

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

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

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

См. также

Следующая