Флаг 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();