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

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

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

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

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

{% note warning %}

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

{% endnote %}

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

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

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

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

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

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

{% list tabs %}

- Java

  ```java
  public class HitParsingFunction implements RowFunction {
      @Override
      public void onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx) {
          var hit = ProtoUtils.parseBytes(message.get("data", byte[].class), THit.class);
          // Дубликаты эмитятся, но не публикуются дальше.
          var distribute = !hit.getHitPayload().equals("duplicate_payload");
          output.addMessage(
                  ctx.createMessageBuilder("hit")
                          .set("hit_id", hit.getHitId())
                          .set("hit_payload", hit.getHitPayload())
                          .finish(),
                  AddMessageOptions.builder().setDistribute(distribute).build()
          );
      }
  }
  ```

- Kotlin

  ```kotlin
  class HitParsingFunction : RowFunction {
      override fun onMessage(message: ExtendedMessage, output: OutputCollector, ctx: RuntimeContext) {
          val hit = ProtoUtils.parseBytes(message.get("data", ByteArray::class.java), THit::class.java)
          // Дубликаты эмитятся, но не публикуются дальше.
          val distribute = hit.hitPayload != "duplicate_payload"
          output.addMessage(
              ctx.createMessageBuilder("hit")
                  .set("hit_id", hit.hitId)
                  .set("hit_payload", hit.hitPayload)
                  .finish(),
              AddMessageOptions.builder().setDistribute(distribute).build()
          )
      }
  }
  ```

{% endlist %}

## Регистрация source-компьютейшена {#registration}

Source-компьютейшен создаётся через `SourceComputation.builder()`. Отдельный параметр фильтрации больше не требуется — решение о публикации принимается в функции обработки.

```java
var hitReader = SourceComputation.builder()
        .setComputationId("hit_reader")
        .setProcessFunction(new HitParsingFunction())
        .build();
```

## См. также

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