JdbcIO.Write.withResults e Wait.on con una PCollection illimitata con FixedWindow
Sto cercando di implementare una pipeline simile a quella delineata in questa domanda , ma a differenza della situazione menzionata in BEAM-6732 , la mia fonte è un abbonamento Pub / Sub e invece di usare il Wait.onper scrivere su un'altra tabella, sono cercando di usarlo per determinare quando le scritture sono complete, generare un messaggio e instradare a un argomento Pub / Sub.
Ho provato a utilizzare la finestra predefinita, ma in base alla documentazione per Wait.on, non funziona per raccolte illimitate, ho provato a definire manualmente una finestra fissa, con un ritardo consentito inferiore, ma anche questo sembra non funzionare, trova la finestra utilizzata di seguito . I passaggi dopo JDBCIO.write sembrano essere sempre bloccati, ovvero non vi è alcun output dal passaggio di attesa.
Window.into(FixedWindows.of(Duration.standardSeconds(10)))
.triggering(
Repeatedly.forever(
AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.standardMinutes(1))
.orFinally(AfterWatermark.pastEndOfWindow())
)
).withAllowedLateness(Duration.standardMinutes(2)).discardingFiredPanes();
Alla ricerca di consigli su cosa potrebbe essere sbagliato, anche su quale sarebbe l'impatto dell'utilizzo di un basso allowedLatenessper una fonte Pub / Sub, che non garantisce l'ordinazione.
Risposte
Sembra che Wait possa essere utilizzato per origini illimitate poiché può essere applicato su Windows. Nella documentazione di Wait SDK troviamo:
"restituisce una PCollection con contenuti identici all'ingresso, ma ritarda la produzione di elementi dell'output nella finestra W fino a quando la finestra del segnale W si chiude (il watermark del segnale supera W.end + signal.allowedLateness)"
Diamo un'occhiata alla spiegazione "apply T to X after Y is ready" -> X.apply(Wait.on(Y)).apply(T). Se ho capito bene, possiamo adattare quanto segue al tuo caso d'uso:
- T = Genera un messaggio
- Y = JdbcIO.Write.withResults
- X = Emit to PubSub
Mentre FixedWindow potrebbe essere richiesto nel tuo frammento di codice, penso che l'attivazione non sarà necessaria poiché Wait sta già considerando la fine della finestra (la filigrana del segnale supera W.end + signal.allowedLateness). Inoltre, con un valore basso per allowedLateness, la finestra verrà chiusa prima, non appena termina la finestra.