Флаг distribute в YTsaurus Flow (Python)
Флаг distribute — это per-message-флаг, задаваемый при добавлении выходного сообщения в SourceComputation. Он управляет тем, будет ли сообщение опубликовано дальше по графу обработки.
Флаг distribute обеспечивает:
- Корректную оценку watermark: сообщения с
distribute=Falseвсё равно учитываются генератором watermark (в отличие от фильтрации вon_message, которая может нарушить watermark). - Присвоение детерминированных идентификаторов сообщениям.
Важно
Чтобы отфильтровать сообщение в SourceComputation, не пропускайте его в on_message — вместо этого эмитьте его с distribute=False. Так сообщение не будет опубликовано дальше, но останется учтённым при оценке watermark.
Когда использовать distribute=False
Флаг distribute=False следует использовать, когда:
- Необходимо отфильтровать часть выходных сообщений на этапе source-компьютейшена.
- Важна корректная оценка watermark.
Если флаг не задан, он по умолчанию равен True, и сообщение публикуется дальше.
Использование
Логика фильтрации переносится в функцию обработки: вместо отдельного шага фильтрации сообщение эмитится с нужным флагом.
from yt.yt.flow.library.python.companion.computation import AddMessageOptions, RowFunction
class HitParsingFunction(RowFunction):
def on_message(self, message, output, ctx):
builder = ctx.message_builder("hit")
builder.set("hit_id", message.payload["hit_id"])
builder.set("hit_payload", message.payload["hit_payload"])
# Дубликаты эмитятся, но не публикуются дальше.
is_duplicate = message.payload["hit_payload"] == "duplicate_payload"
output.add_message(
builder.finish(),
AddMessageOptions(distribute=not is_duplicate),
)
Регистрация source-компьютейшена
Source-компьютейшен регистрируется через Pipeline.add() с source=True. Отдельный параметр фильтрации больше не требуется — решение о публикации принимается в функции обработки.
from yt.yt.flow.library.python.companion import Pipeline
pipeline = Pipeline()
pipeline.add("hit_reader", HitParsingFunction(), source=True)