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

<!-- source: ru/_includes/flow/python/distribute.md -->
# Флаг distribute в YTsaurus Flow (Python)

Флаг `distribute` — это per-message-флаг, задаваемый при добавлении выходного [сообщения](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#message) в [SourceComputation](https://ytsaurus.tech/docs/ru/flow/python/computation.md#sourcecomputation). Он управляет тем, будет ли сообщение опубликовано дальше по графу обработки.

Флаг `distribute` обеспечивает:

- Корректную оценку [watermark](https://ytsaurus.tech/docs/ru/flow/concepts/watermarks.md): сообщения с `distribute=False` всё равно учитываются генератором watermark (в отличие от фильтрации в `on_message`, которая может нарушить watermark).
- Присвоение детерминированных идентификаторов сообщениям.

{% note warning %}

Чтобы отфильтровать сообщение в `SourceComputation`, не пропускайте его в `on_message` — вместо этого эмитьте его с `distribute=False`. Так сообщение не будет опубликовано дальше, но останется учтённым при оценке watermark.

{% endnote %}

## Когда использовать distribute=False

Флаг `distribute=False` следует использовать, когда:

- Необходимо отфильтровать часть выходных сообщений на этапе source-компьютейшена.
- Важна корректная оценка watermark.

Если флаг не задан, он по умолчанию равен `True`, и сообщение публикуется дальше.

## Использование {#usage}

Логика фильтрации переносится в функцию обработки: вместо отдельного шага фильтрации сообщение эмитится с нужным флагом.

```python
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-компьютейшена {#registration}

Source-компьютейшен регистрируется через `Pipeline.add()` с `source=True`. Отдельный параметр фильтрации больше не требуется — решение о публикации принимается в функции обработки.

```python
from yt.yt.flow.library.python.companion import Pipeline

pipeline = Pipeline()
pipeline.add("hit_reader", HitParsingFunction(), source=True)
```

## См. также

- [Computation (Python)](https://ytsaurus.tech/docs/ru/flow/python/computation.md)
- [Watermarks](https://ytsaurus.tech/docs/ru/flow/concepts/watermarks.md)
- [Флаг distribute (Java)](https://ytsaurus.tech/docs/ru/flow/java/distribute.md)
<!-- endsource: ru/_includes/flow/python/distribute.md -->
