JdbcIO.Write.withResults e Wait.on com uma PCollection ilimitada com FixedWindow
Estou tentando implementar um pipeline semelhante ao descrito nesta pergunta , mas, ao contrário da situação mencionada em BEAM-6732 , minha fonte é uma assinatura Pub / Sub e, em vez de usar o Wait.onpara gravar em outra tabela, estou tentar usá-lo para determinar quando as gravações estão concluídas, gerar uma mensagem e encaminhar para um tópico Pub / Sub.
Tentei usar a janela padrão, mas com base na documentação de Wait.on, ela não funciona para coleções ilimitadas, tentei definir manualmente uma janela fixa, com um atraso inferior permitido, mas que também parece não funcionar, encontre a janela usada abaixo . As etapas após o JDBCIO.write parecem estar sempre travadas, ou seja, não há saída da etapa de espera.
Window.into(FixedWindows.of(Duration.standardSeconds(10)))
.triggering(
Repeatedly.forever(
AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.standardMinutes(1))
.orFinally(AfterWatermark.pastEndOfWindow())
)
).withAllowedLateness(Duration.standardMinutes(2)).discardingFiredPanes();
Em busca de conselhos sobre o que pode estar errado, também qual seria o impacto de usar um baixo allowedLatenesspara uma fonte Pub / Sub, o que não garante o pedido.
Respostas
Parece que Wait pode ser usado para fontes ilimitadas, pois pode ser aplicado no Windows. Na documentação do Wait SDK , encontramos:
"retorna uma PCollection com conteúdo idêntico ao da entrada, mas atrasa a produção de elementos da saída na janela W até que a janela W do sinal feche (a marca d'água do sinal passa de W.end + signal.allowedLateness)"
Vejamos a explicação "apply T to X after Y is ready" -> X.apply(Wait.on(Y)).apply(T). Se bem entendi, podemos adaptar o seguinte ao seu caso de uso:
- T = Gerar uma mensagem
- Y = JdbcIO.Write.withResults
- X = emitir para PubSub
Embora FixedWindow possa ser necessário em seu trecho de código, acho que o acionamento não será necessário, pois Wait já está considerando o fim da janela (a marca d'água do sinal passa W.end + signal.allowedLateness). Além disso, com um valor baixo para allowedLateness, a janela será fechada mais cedo, assim que a janela terminar.