Порядок обработки сообщений в 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 среди сообщений, находящихся в данный момент во всех выходных буферах по данному стриму:
Такой подход позволяет учитывать специфику каждого стрима (например, систематическую задержку между временем записи в очередь и временем события) и обеспечивает более справедливое межстримовое упорядочивание.