Флаг distribute в YTsaurus Flow (Java)

Флаг distribute — это per-message-флаг, задаваемый при добавлении выходного сообщения в SourceComputation. Он управляет тем, будет ли сообщение опубликовано дальше по графу обработки.

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

  • Корректную оценку watermark: сообщения с distribute=false всё равно учитываются генератором watermark (в отличие от фильтрации в onMessage, которая может нарушить watermark).
  • Присвоение детерминированных идентификаторов сообщениям.

Важно

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

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

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

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

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

Использование

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

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()
        );
    }
}
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()
        )
    }
}

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

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

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

См. также

Предыдущая
Следующая